From 1ccd61ece5ecf57c8be3f3f36d62cd500943b74f Mon Sep 17 00:00:00 2001 From: fogelito Date: Tue, 3 Mar 2026 12:15:21 +0200 Subject: [PATCH 01/13] functionInternalId --- app/config/collections/projects.php | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/app/config/collections/projects.php b/app/config/collections/projects.php index ecd3db1733..107e1061b9 100644 --- a/app/config/collections/projects.php +++ b/app/config/collections/projects.php @@ -2193,7 +2193,7 @@ return [ [ '$id' => ID::custom('_key_function_internal_id'), 'type' => Database::INDEX_KEY, - 'attributes' => ['resourceInternalId'], + 'attributes' => ['functionInternalId'], 'lengths' => [], 'orders' => [], ], From f4125b885956e6eecc817772508bfe39627488ef Mon Sep 17 00:00:00 2001 From: fogelito Date: Tue, 3 Mar 2026 13:35:23 +0200 Subject: [PATCH 02/13] Remove length --- app/config/collections/projects.php | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/app/config/collections/projects.php b/app/config/collections/projects.php index 107e1061b9..b84acf9628 100644 --- a/app/config/collections/projects.php +++ b/app/config/collections/projects.php @@ -2145,21 +2145,21 @@ return [ '$id' => ID::custom('_key_trigger'), 'type' => Database::INDEX_KEY, 'attributes' => ['trigger'], - 'lengths' => [32], + 'lengths' => [], 'orders' => [Database::ORDER_ASC], ], [ '$id' => ID::custom('_key_status'), 'type' => Database::INDEX_KEY, 'attributes' => ['status'], - 'lengths' => [32], + 'lengths' => [], 'orders' => [Database::ORDER_ASC], ], [ '$id' => ID::custom('_key_requestMethod'), 'type' => Database::INDEX_KEY, 'attributes' => ['requestMethod'], - 'lengths' => [128], + 'lengths' => [], 'orders' => [Database::ORDER_ASC], ], [ From b6b44efdab0709e0317125effebd1ed81c3538d6 Mon Sep 17 00:00:00 2001 From: fogelito Date: Tue, 3 Mar 2026 13:58:00 +0200 Subject: [PATCH 03/13] Revert length --- app/config/collections/projects.php | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/app/config/collections/projects.php b/app/config/collections/projects.php index b84acf9628..107e1061b9 100644 --- a/app/config/collections/projects.php +++ b/app/config/collections/projects.php @@ -2145,21 +2145,21 @@ return [ '$id' => ID::custom('_key_trigger'), 'type' => Database::INDEX_KEY, 'attributes' => ['trigger'], - 'lengths' => [], + 'lengths' => [32], 'orders' => [Database::ORDER_ASC], ], [ '$id' => ID::custom('_key_status'), 'type' => Database::INDEX_KEY, 'attributes' => ['status'], - 'lengths' => [], + 'lengths' => [32], 'orders' => [Database::ORDER_ASC], ], [ '$id' => ID::custom('_key_requestMethod'), 'type' => Database::INDEX_KEY, 'attributes' => ['requestMethod'], - 'lengths' => [], + 'lengths' => [128], 'orders' => [Database::ORDER_ASC], ], [ From 5999db295dd696652146e257047de4d54eca1eaa Mon Sep 17 00:00:00 2001 From: fogelito Date: Tue, 3 Mar 2026 15:44:20 +0200 Subject: [PATCH 04/13] Message --- app/config/collections/projects.php | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/app/config/collections/projects.php b/app/config/collections/projects.php index 107e1061b9..da37334b4e 100644 --- a/app/config/collections/projects.php +++ b/app/config/collections/projects.php @@ -2191,7 +2191,7 @@ return [ 'orders' => [Database::ORDER_ASC], ], [ - '$id' => ID::custom('_key_function_internal_id'), + '$id' => ID::custom('_key_function_internal_id'), // Index not in use remove in the future 'type' => Database::INDEX_KEY, 'attributes' => ['functionInternalId'], 'lengths' => [], From b45ff6b6461219677faed24dde444caa9f50ed26 Mon Sep 17 00:00:00 2001 From: loks0n <22452787+loks0n@users.noreply.github.com> Date: Fri, 13 Feb 2026 12:28:26 +0000 Subject: [PATCH 05/13] refactor: replace queueForExecutions with Bus event bus Introduce a generic event bus (Utopia\Bus) with typed events, listener base class, Span instrumentation, and coroutine dispatch. Replace all direct queueForExecutions and inline execution usage calls with ExecutionCompleted event and dedicated listeners (Log, Usage). Co-Authored-By: Claude Opus 4.6 --- app/cli.php | 4 ++ app/controllers/general.php | 73 +++++++-------------- app/http.php | 4 ++ app/init/registers.php | 8 +++ app/init/resources.php | 4 -- app/listeners.php | 9 +++ app/worker.php | 8 +-- composer.json | 3 +- phpstan.neon | 4 ++ src/Appwrite/Bus/ExecutionCompleted.php | 22 +++++++ src/Appwrite/Bus/Listeners/Log.php | 39 +++++++++++ src/Appwrite/Bus/Listeners/Usage.php | 64 ++++++++++++++++++ src/Appwrite/Platform/Workers/Functions.php | 65 +++++++----------- src/Utopia/Bus/Bus.php | 59 +++++++++++++++++ src/Utopia/Bus/Event.php | 7 ++ src/Utopia/Bus/Listener.php | 51 ++++++++++++++ 16 files changed, 321 insertions(+), 103 deletions(-) create mode 100644 app/listeners.php create mode 100644 src/Appwrite/Bus/ExecutionCompleted.php create mode 100644 src/Appwrite/Bus/Listeners/Log.php create mode 100644 src/Appwrite/Bus/Listeners/Usage.php create mode 100644 src/Utopia/Bus/Bus.php create mode 100644 src/Utopia/Bus/Event.php create mode 100644 src/Utopia/Bus/Listener.php diff --git a/app/cli.php b/app/cli.php index 0f8426afd9..052643f004 100644 --- a/app/cli.php +++ b/app/cli.php @@ -318,6 +318,10 @@ $setResource('logError', function (Registry $register) { $setResource('executor', fn () => new Executor(), []); +$setResource('bus', function (Registry $register) use ($cli) { + return $register->get('bus')->setResolver(fn (string $name) => $cli->getResource($name)); +}, ['register']); + $setResource('telemetry', fn () => new NoTelemetry(), []); $cli diff --git a/app/controllers/general.php b/app/controllers/general.php index 2ac03368df..5f1a5e5d9d 100644 --- a/app/controllers/general.php +++ b/app/controllers/general.php @@ -5,10 +5,10 @@ require_once __DIR__ . '/../init.php'; use Ahc\Jwt\JWT; use Ahc\Jwt\JWTException; use Appwrite\Auth\Key; +use Appwrite\Bus\ExecutionCompleted; use Appwrite\Event\Certificate; use Appwrite\Event\Delete as DeleteEvent; use Appwrite\Event\Event; -use Appwrite\Event\Execution; use Appwrite\Event\StatsUsage; use Appwrite\Extend\Exception as AppwriteException; use Appwrite\Network\Cors; @@ -35,6 +35,7 @@ use Executor\Executor; use MaxMind\Db\Reader; use Swoole\Http\Request as SwooleRequest; use Swoole\Table; +use Utopia\Bus\Bus; use Utopia\Config\Config; use Utopia\Console; use Utopia\Database\Database; @@ -62,7 +63,7 @@ Config::setParam('domainVerification', false); Config::setParam('cookieDomain', 'localhost'); Config::setParam('cookieSamesite', Response::COOKIE_SAMESITE_NONE); -function router(Http $utopia, Database $dbForPlatform, callable $getProjectDB, SwooleRequest $swooleRequest, Request $request, Response $response, Log $log, Event $queueForEvents, StatsUsage $queueForStatsUsage, Execution $queueForExecutions, Executor $executor, Reader $geodb, callable $isResourceBlocked, array $platform, string $previewHostname, Authorization $authorization, ?Key $apiKey, DeleteEvent $queueForDeletes, int $executionsRetentionCount) +function router(Http $utopia, Database $dbForPlatform, callable $getProjectDB, SwooleRequest $swooleRequest, Request $request, Response $response, Log $log, Event $queueForEvents, StatsUsage $queueForStatsUsage, Bus $bus, Executor $executor, Reader $geodb, callable $isResourceBlocked, array $platform, string $previewHostname, Authorization $authorization, ?Key $apiKey, DeleteEvent $queueForDeletes, int $executionsRetentionCount) { $host = $request->getHostname() ?? ''; if (!empty($previewHostname)) { @@ -706,10 +707,12 @@ function router(Http $utopia, Database $dbForPlatform, callable $getProjectDB, S } } finally { if ($type === 'function' || $type === 'site') { - $queueForExecutions - ->setExecution($execution) - ->setProject($project) - ->trigger(); + $bus->dispatch(new ExecutionCompleted( + execution: $execution->getArrayCopy(), + project: $project->getArrayCopy(), + spec: $spec, + resource: $resource->getArrayCopy(), + )); } } @@ -766,31 +769,11 @@ function router(Http $utopia, Database $dbForPlatform, callable $getProjectDB, S } } - $metricTypeExecutions = str_replace(['{resourceType}'], [$deployment->getAttribute('resourceType')], METRIC_RESOURCE_TYPE_EXECUTIONS); - $metricTypeIdExecutions = str_replace(['{resourceType}', '{resourceInternalId}'], [$deployment->getAttribute('resourceType'), $resource->getSequence()], METRIC_RESOURCE_TYPE_ID_EXECUTIONS); - $metricTypeExecutionsCompute = str_replace(['{resourceType}'], [$deployment->getAttribute('resourceType')], METRIC_RESOURCE_TYPE_EXECUTIONS_COMPUTE); - $metricTypeIdExecutionsCompute = str_replace(['{resourceType}', '{resourceInternalId}'], [$deployment->getAttribute('resourceType'), $resource->getSequence()], METRIC_RESOURCE_TYPE_ID_EXECUTIONS_COMPUTE); - $metricTypeExecutionsMbSeconds = str_replace(['{resourceType}'], [$deployment->getAttribute('resourceType')], METRIC_RESOURCE_TYPE_EXECUTIONS_MB_SECONDS); - $metricTypeIdExecutionsMBSeconds = str_replace(['{resourceType}', '{resourceInternalId}'], [$deployment->getAttribute('resourceType'), $resource->getSequence()], METRIC_RESOURCE_TYPE_ID_EXECUTIONS_MB_SECONDS); if ($deployment->getAttribute('resourceType') === 'sites') { $queueForStatsUsage ->disableMetric(METRIC_NETWORK_REQUESTS) ->disableMetric(METRIC_NETWORK_INBOUND) - ->disableMetric(METRIC_NETWORK_OUTBOUND); - if ($resource->getAttribute('adapter') !== 'ssr') { - $queueForStatsUsage - ->disableMetric(METRIC_EXECUTIONS) - ->disableMetric(METRIC_EXECUTIONS_COMPUTE) - ->disableMetric(METRIC_EXECUTIONS_MB_SECONDS) - ->disableMetric($metricTypeExecutions) - ->disableMetric($metricTypeIdExecutions) - ->disableMetric($metricTypeExecutionsCompute) - ->disableMetric($metricTypeIdExecutionsCompute) - ->disableMetric($metricTypeExecutionsMbSeconds) - ->disableMetric($metricTypeIdExecutionsMBSeconds); - } - - $queueForStatsUsage + ->disableMetric(METRIC_NETWORK_OUTBOUND) ->addMetric(METRIC_SITES_REQUESTS, 1) ->addMetric(METRIC_SITES_INBOUND, $request->getSize() + $fileSize) ->addMetric(METRIC_SITES_OUTBOUND, $response->getSize()) @@ -800,22 +783,10 @@ function router(Http $utopia, Database $dbForPlatform, callable $getProjectDB, S ; } - - $compute = (int)($execution->getAttribute('duration') * 1000); - $mbSeconds = (int)(($spec['memory'] ?? APP_COMPUTE_MEMORY_DEFAULT) * $execution->getAttribute('duration', 0) * ($spec['cpus'] ?? APP_COMPUTE_CPUS_DEFAULT)); $queueForStatsUsage ->addMetric(METRIC_NETWORK_REQUESTS, 1) ->addMetric(METRIC_NETWORK_INBOUND, $request->getSize() + $fileSize) ->addMetric(METRIC_NETWORK_OUTBOUND, $response->getSize()) - ->addMetric(METRIC_EXECUTIONS, 1) - ->addMetric($metricTypeExecutions, 1) - ->addMetric($metricTypeIdExecutions, 1) - ->addMetric(METRIC_EXECUTIONS_COMPUTE, $compute) // per project - ->addMetric($metricTypeExecutionsCompute, $compute) // per function - ->addMetric($metricTypeIdExecutionsCompute, $compute) // per function - ->addMetric(METRIC_EXECUTIONS_MB_SECONDS, $mbSeconds) - ->addMetric($metricTypeExecutionsMbSeconds, $mbSeconds) - ->addMetric($metricTypeIdExecutionsMBSeconds, $mbSeconds) ->setProject($project) ->trigger(); @@ -883,7 +854,7 @@ Http::init() ->inject('geodb') ->inject('queueForStatsUsage') ->inject('queueForEvents') - ->inject('queueForExecutions') + ->inject('bus') ->inject('executor') ->inject('platform') ->inject('isResourceBlocked') @@ -894,7 +865,7 @@ Http::init() ->inject('authorization') ->inject('queueForDeletes') ->inject('executionsRetentionCount') - ->action(function (Http $utopia, SwooleRequest $swooleRequest, Request $request, Response $response, Log $log, Document $project, Database $dbForPlatform, callable $getProjectDB, Locale $locale, array $localeCodes, Reader $geodb, StatsUsage $queueForStatsUsage, Event $queueForEvents, Execution $queueForExecutions, Executor $executor, array $platform, callable $isResourceBlocked, string $previewHostname, Document $devKey, ?Key $apiKey, Cors $cors, Authorization $authorization, DeleteEvent $queueForDeletes, int $executionsRetentionCount) { + ->action(function (Http $utopia, SwooleRequest $swooleRequest, Request $request, Response $response, Log $log, Document $project, Database $dbForPlatform, callable $getProjectDB, Locale $locale, array $localeCodes, Reader $geodb, StatsUsage $queueForStatsUsage, Event $queueForEvents, Bus $bus, Executor $executor, array $platform, callable $isResourceBlocked, string $previewHostname, Document $devKey, ?Key $apiKey, Cors $cors, Authorization $authorization, DeleteEvent $queueForDeletes, int $executionsRetentionCount) { /* * Appwrite Router */ @@ -902,7 +873,7 @@ Http::init() $platformHostnames = $platform['hostnames'] ?? []; // Only run Router when external domain if (!\in_array($hostname, $platformHostnames) || !empty($previewHostname)) { - if (router($utopia, $dbForPlatform, $getProjectDB, $swooleRequest, $request, $response, $log, $queueForEvents, $queueForStatsUsage, $queueForExecutions, $executor, $geodb, $isResourceBlocked, $platform, $previewHostname, $authorization, $apiKey, $queueForDeletes, $executionsRetentionCount)) { + if (router($utopia, $dbForPlatform, $getProjectDB, $swooleRequest, $request, $response, $log, $queueForEvents, $queueForStatsUsage, $bus, $executor, $geodb, $isResourceBlocked, $platform, $previewHostname, $authorization, $apiKey, $queueForDeletes, $executionsRetentionCount)) { $utopia->getRoute()?->label('router', true); } } @@ -1179,7 +1150,7 @@ Http::options() ->inject('getProjectDB') ->inject('queueForEvents') ->inject('queueForStatsUsage') - ->inject('queueForExecutions') + ->inject('bus') ->inject('executor') ->inject('geodb') ->inject('isResourceBlocked') @@ -1192,14 +1163,14 @@ Http::options() ->inject('authorization') ->inject('queueForDeletes') ->inject('executionsRetentionCount') - ->action(function (Http $utopia, SwooleRequest $swooleRequest, Request $request, Response $response, Log $log, Database $dbForPlatform, callable $getProjectDB, Event $queueForEvents, StatsUsage $queueForStatsUsage, Execution $queueForExecutions, Executor $executor, Reader $geodb, callable $isResourceBlocked, array $platform, string $previewHostname, Document $project, Document $devKey, ?Key $apiKey, Cors $cors, Authorization $authorization, DeleteEvent $queueForDeletes, int $executionsRetentionCount) { + ->action(function (Http $utopia, SwooleRequest $swooleRequest, Request $request, Response $response, Log $log, Database $dbForPlatform, callable $getProjectDB, Event $queueForEvents, StatsUsage $queueForStatsUsage, Bus $bus, Executor $executor, Reader $geodb, callable $isResourceBlocked, array $platform, string $previewHostname, Document $project, Document $devKey, ?Key $apiKey, Cors $cors, Authorization $authorization, DeleteEvent $queueForDeletes, int $executionsRetentionCount) { /* * Appwrite Router */ $platformHostnames = $platform['hostnames'] ?? []; // Only run Router when external domain if (!in_array($request->getHostname(), $platformHostnames) || !empty($previewHostname)) { - if (router($utopia, $dbForPlatform, $getProjectDB, $swooleRequest, $request, $response, $log, $queueForEvents, $queueForStatsUsage, $queueForExecutions, $executor, $geodb, $isResourceBlocked, $platform, $previewHostname, $authorization, $apiKey, $queueForDeletes, $executionsRetentionCount)) { + if (router($utopia, $dbForPlatform, $getProjectDB, $swooleRequest, $request, $response, $log, $queueForEvents, $queueForStatsUsage, $bus, $executor, $geodb, $isResourceBlocked, $platform, $previewHostname, $authorization, $apiKey, $queueForDeletes, $executionsRetentionCount)) { $utopia->getRoute()?->label('router', true); } } @@ -1569,7 +1540,7 @@ Http::get('/robots.txt') ->inject('getProjectDB') ->inject('queueForEvents') ->inject('queueForStatsUsage') - ->inject('queueForExecutions') + ->inject('bus') ->inject('executor') ->inject('geodb') ->inject('isResourceBlocked') @@ -1579,13 +1550,13 @@ Http::get('/robots.txt') ->inject('authorization') ->inject('queueForDeletes') ->inject('executionsRetentionCount') - ->action(function (Http $utopia, SwooleRequest $swooleRequest, Request $request, Response $response, Log $log, Database $dbForPlatform, callable $getProjectDB, Event $queueForEvents, StatsUsage $queueForStatsUsage, Execution $queueForExecutions, Executor $executor, Reader $geodb, callable $isResourceBlocked, array $platform, string $previewHostname, ?Key $apiKey, Authorization $authorization, DeleteEvent $queueForDeletes, int $executionsRetentionCount) { + ->action(function (Http $utopia, SwooleRequest $swooleRequest, Request $request, Response $response, Log $log, Database $dbForPlatform, callable $getProjectDB, Event $queueForEvents, StatsUsage $queueForStatsUsage, Bus $bus, Executor $executor, Reader $geodb, callable $isResourceBlocked, array $platform, string $previewHostname, ?Key $apiKey, Authorization $authorization, DeleteEvent $queueForDeletes, int $executionsRetentionCount) { $platformHostnames = $platform['hostnames'] ?? []; if (in_array($request->getHostname(), $platformHostnames) || !empty($previewHostname)) { $template = new View(__DIR__ . '/../views/general/robots.phtml'); $response->text($template->render(false)); } else { - if (router($utopia, $dbForPlatform, $getProjectDB, $swooleRequest, $request, $response, $log, $queueForEvents, $queueForStatsUsage, $queueForExecutions, $executor, $geodb, $isResourceBlocked, $platform, $previewHostname, $authorization, $apiKey, $queueForDeletes, $executionsRetentionCount)) { + if (router($utopia, $dbForPlatform, $getProjectDB, $swooleRequest, $request, $response, $log, $queueForEvents, $queueForStatsUsage, $bus, $executor, $geodb, $isResourceBlocked, $platform, $previewHostname, $authorization, $apiKey, $queueForDeletes, $executionsRetentionCount)) { $utopia->getRoute()?->label('router', true); } } @@ -1604,7 +1575,7 @@ Http::get('/humans.txt') ->inject('getProjectDB') ->inject('queueForEvents') ->inject('queueForStatsUsage') - ->inject('queueForExecutions') + ->inject('bus') ->inject('executor') ->inject('geodb') ->inject('isResourceBlocked') @@ -1614,13 +1585,13 @@ Http::get('/humans.txt') ->inject('authorization') ->inject('queueForDeletes') ->inject('executionsRetentionCount') - ->action(function (Http $utopia, SwooleRequest $swooleRequest, Request $request, Response $response, Log $log, Database $dbForPlatform, callable $getProjectDB, Event $queueForEvents, StatsUsage $queueForStatsUsage, Execution $queueForExecutions, Executor $executor, Reader $geodb, callable $isResourceBlocked, array $platform, string $previewHostname, ?Key $apiKey, Authorization $authorization, DeleteEvent $queueForDeletes, int $executionsRetentionCount) { + ->action(function (Http $utopia, SwooleRequest $swooleRequest, Request $request, Response $response, Log $log, Database $dbForPlatform, callable $getProjectDB, Event $queueForEvents, StatsUsage $queueForStatsUsage, Bus $bus, Executor $executor, Reader $geodb, callable $isResourceBlocked, array $platform, string $previewHostname, ?Key $apiKey, Authorization $authorization, DeleteEvent $queueForDeletes, int $executionsRetentionCount) { $platformHostnames = $platform['hostnames'] ?? []; if (in_array($request->getHostname(), $platformHostnames) || !empty($previewHostname)) { $template = new View(__DIR__ . '/../views/general/humans.phtml'); $response->text($template->render(false)); } else { - if (router($utopia, $dbForPlatform, $getProjectDB, $swooleRequest, $request, $response, $log, $queueForEvents, $queueForStatsUsage, $queueForExecutions, $executor, $geodb, $isResourceBlocked, $platform, $previewHostname, $authorization, $apiKey, $queueForDeletes, $executionsRetentionCount)) { + if (router($utopia, $dbForPlatform, $getProjectDB, $swooleRequest, $request, $response, $log, $queueForEvents, $queueForStatsUsage, $bus, $executor, $geodb, $isResourceBlocked, $platform, $previewHostname, $authorization, $apiKey, $queueForDeletes, $executionsRetentionCount)) { $utopia->getRoute()?->label('router', true); } } diff --git a/app/http.php b/app/http.php index 30c4a8317b..7f771de130 100644 --- a/app/http.php +++ b/app/http.php @@ -188,6 +188,10 @@ $http->on(Constant::EVENT_AFTER_RELOAD, function ($server) { Console::success('Reload completed...'); }); +Http::setResource('bus', function ($register, $utopia) { + return $register->get('bus')->setResolver(fn (string $name) => $utopia->getResource($name)); +}, ['register', 'utopia']); + include __DIR__ . '/controllers/general.php'; function createDatabase(Http $app, string $resourceKey, string $dbName, array $collections, mixed $pools, ?callable $extraSetup = null): void diff --git a/app/init/registers.php b/app/init/registers.php index 411fd4c69d..26a9329270 100644 --- a/app/init/registers.php +++ b/app/init/registers.php @@ -449,3 +449,11 @@ $register->set('promiseAdapter', function () { $register->set('hooks', function () { return new Hooks(); }); +$listeners = require __DIR__ . '/../listeners.php'; +$register->set('bus', function () use ($listeners) { + $bus = new \Utopia\Bus\Bus(); + foreach ($listeners as $listener) { + $bus->subscribe($listener); + } + return $bus; +}); diff --git a/app/init/resources.php b/app/init/resources.php index f3e66c388a..36aadc9707 100644 --- a/app/init/resources.php +++ b/app/init/resources.php @@ -10,7 +10,6 @@ use Appwrite\Event\Certificate; use Appwrite\Event\Database as EventDatabase; use Appwrite\Event\Delete; use Appwrite\Event\Event; -use Appwrite\Event\Execution; use Appwrite\Event\Func; use Appwrite\Event\Mail; use Appwrite\Event\Messaging; @@ -160,9 +159,6 @@ Http::setResource('queueForAudits', function (Publisher $publisher) { Http::setResource('queueForFunctions', function (Publisher $publisher) { return new Func($publisher); }, ['publisher']); -Http::setResource('queueForExecutions', function (Publisher $publisher) { - return new Execution($publisher); -}, ['publisher']); Http::setResource('eventProcessor', function () { return new EventProcessor(); }, []); diff --git a/app/listeners.php b/app/listeners.php new file mode 100644 index 0000000000..714c255974 --- /dev/null +++ b/app/listeners.php @@ -0,0 +1,9 @@ +getWorker(); +Server::setResource('bus', function ($register) use ($worker) { + return $register->get('bus')->setResolver(fn (string $name) => $worker->getResource($name)); +}, ['register']); + $worker ->error() ->inject('error') diff --git a/composer.json b/composer.json index 3dd24aa29f..1422bd5d0a 100644 --- a/composer.json +++ b/composer.json @@ -19,7 +19,8 @@ "autoload": { "psr-4": { "Appwrite\\": "src/Appwrite", - "Executor\\": "src/Executor" + "Executor\\": "src/Executor", + "Utopia\\Bus\\": "src/Utopia/Bus" } }, "autoload-dev": { diff --git a/phpstan.neon b/phpstan.neon index 153b3be21c..90f28e7539 100644 --- a/phpstan.neon +++ b/phpstan.neon @@ -1,7 +1,11 @@ parameters: level: 8 paths: + - src/Utopia/Bus + - src/Appwrite/Bus - src/Appwrite/Transformation + bootstrapFiles: + - app/init/constants.php scanDirectories: - vendor/swoole/ide-helper excludePaths: diff --git a/src/Appwrite/Bus/ExecutionCompleted.php b/src/Appwrite/Bus/ExecutionCompleted.php new file mode 100644 index 0000000000..1508266ea0 --- /dev/null +++ b/src/Appwrite/Bus/ExecutionCompleted.php @@ -0,0 +1,22 @@ + $execution + * @param array $project + * @param array $spec + * @param array $resource + */ + public function __construct( + public readonly array $execution, + public readonly array $project, + public readonly array $spec = [], + public readonly array $resource = [], + ) { + } +} diff --git a/src/Appwrite/Bus/Listeners/Log.php b/src/Appwrite/Bus/Listeners/Log.php new file mode 100644 index 0000000000..12ca7303c8 --- /dev/null +++ b/src/Appwrite/Bus/Listeners/Log.php @@ -0,0 +1,39 @@ +desc('Persists execution logs to database via queue') + ->inject('publisher') + ->callback($this->handle(...)); + } + + public function handle(ExecutionCompleted $event, Publisher $publisher): void + { + $queueForExecutions = new Execution($publisher); + $queueForExecutions + ->setExecution(new Document($event->execution)) + ->setProject(new Document($event->project)) + ->trigger(); + } +} diff --git a/src/Appwrite/Bus/Listeners/Usage.php b/src/Appwrite/Bus/Listeners/Usage.php new file mode 100644 index 0000000000..e91bd910f5 --- /dev/null +++ b/src/Appwrite/Bus/Listeners/Usage.php @@ -0,0 +1,64 @@ +desc('Records execution usage metrics') + ->inject('publisher') + ->callback($this->handle(...)); + } + + public function handle(ExecutionCompleted $event, Publisher $publisher): void + { + $execution = new Document($event->execution); + $resource = new Document($event->resource); + + // Non-SSR sites don't record execution metrics + if ($execution->getAttribute('resourceType') === 'sites' && $resource->getAttribute('adapter') !== 'ssr') { + return; + } + $project = new Document($event->project); + $spec = $event->spec; + + $resourceType = $execution->getAttribute('resourceType', ''); + $resourceInternalId = $execution->getAttribute('resourceInternalId', ''); + $duration = $execution->getAttribute('duration', 0); + + $compute = (int)($duration * 1000); + $mbSeconds = (int)(($spec['memory'] ?? APP_COMPUTE_MEMORY_DEFAULT) * $duration * ($spec['cpus'] ?? APP_COMPUTE_CPUS_DEFAULT)); + + $queueForStatsUsage = new StatsUsage($publisher); + $queueForStatsUsage + ->setProject($project) + ->addMetric(METRIC_EXECUTIONS, 1) + ->addMetric(str_replace(['{resourceType}'], [$resourceType], METRIC_RESOURCE_TYPE_EXECUTIONS), 1) + ->addMetric(str_replace(['{resourceType}', '{resourceInternalId}'], [$resourceType, $resourceInternalId], METRIC_RESOURCE_TYPE_ID_EXECUTIONS), 1) + ->addMetric(METRIC_EXECUTIONS_COMPUTE, $compute) + ->addMetric(str_replace(['{resourceType}'], [$resourceType], METRIC_RESOURCE_TYPE_EXECUTIONS_COMPUTE), $compute) + ->addMetric(str_replace(['{resourceType}', '{resourceInternalId}'], [$resourceType, $resourceInternalId], METRIC_RESOURCE_TYPE_ID_EXECUTIONS_COMPUTE), $compute) + ->addMetric(METRIC_EXECUTIONS_MB_SECONDS, $mbSeconds) + ->addMetric(str_replace(['{resourceType}'], [$resourceType], METRIC_RESOURCE_TYPE_EXECUTIONS_MB_SECONDS), $mbSeconds) + ->addMetric(str_replace(['{resourceType}', '{resourceInternalId}'], [$resourceType, $resourceInternalId], METRIC_RESOURCE_TYPE_ID_EXECUTIONS_MB_SECONDS), $mbSeconds) + ->trigger(); + } +} diff --git a/src/Appwrite/Platform/Workers/Functions.php b/src/Appwrite/Platform/Workers/Functions.php index 0932aea335..6c59bc8da0 100644 --- a/src/Appwrite/Platform/Workers/Functions.php +++ b/src/Appwrite/Platform/Workers/Functions.php @@ -3,15 +3,15 @@ namespace Appwrite\Platform\Workers; use Ahc\Jwt\JWT; +use Appwrite\Bus\ExecutionCompleted; use Appwrite\Event\Event; -use Appwrite\Event\Execution as ExecutionEvent; use Appwrite\Event\Func; use Appwrite\Event\Realtime; -use Appwrite\Event\StatsUsage; use Appwrite\Event\Webhook; use Appwrite\Extend\Exception as AppwriteException; use Appwrite\Utopia\Response\Model\Execution; use Executor\Executor; +use Utopia\Bus\Bus; use Utopia\Config\Config; use Utopia\Console; use Utopia\Database\Database; @@ -47,8 +47,7 @@ class Functions extends Action ->inject('queueForFunctions') ->inject('queueForRealtime') ->inject('queueForEvents') - ->inject('queueForStatsUsage') - ->inject('queueForExecutions') + ->inject('bus') ->inject('log') ->inject('executor') ->inject('isResourceBlocked') @@ -63,8 +62,7 @@ class Functions extends Action Func $queueForFunctions, Realtime $queueForRealtime, Event $queueForEvents, - StatsUsage $queueForStatsUsage, - ExecutionEvent $queueForExecutions, + Bus $bus, Log $log, Executor $executor, callable $isResourceBlocked @@ -158,9 +156,8 @@ class Functions extends Action queueForWebhooks: $queueForWebhooks, queueForFunctions: $queueForFunctions, queueForRealtime: $queueForRealtime, - queueForStatsUsage: $queueForStatsUsage, queueForEvents: $queueForEvents, - queueForExecutions: $queueForExecutions, + bus: $bus, project: $project, function: $function, executor: $executor, @@ -203,9 +200,8 @@ class Functions extends Action queueForWebhooks: $queueForWebhooks, queueForFunctions: $queueForFunctions, queueForRealtime: $queueForRealtime, - queueForStatsUsage: $queueForStatsUsage, queueForEvents: $queueForEvents, - queueForExecutions: $queueForExecutions, + bus: $bus, project: $project, function: $function, executor: $executor, @@ -230,9 +226,8 @@ class Functions extends Action queueForWebhooks: $queueForWebhooks, queueForFunctions: $queueForFunctions, queueForRealtime: $queueForRealtime, - queueForStatsUsage: $queueForStatsUsage, queueForEvents: $queueForEvents, - queueForExecutions: $queueForExecutions, + bus: $bus, project: $project, function: $function, executor: $executor, @@ -266,7 +261,7 @@ class Functions extends Action private function fail( string $message, Document $project, - ExecutionEvent $queueForExecutions, + Bus $bus, Document $function, string $trigger, string $path, @@ -309,10 +304,10 @@ class Functions extends Action 'duration' => 0.0, ]); - $queueForExecutions - ->setExecution($execution) - ->setProject($project) - ->trigger(); + $bus->dispatch(new ExecutionCompleted( + execution: $execution->getArrayCopy(), + project: $project->getArrayCopy(), + )); } /** @@ -320,7 +315,6 @@ class Functions extends Action * @param Database $dbForProject * @param Func $queueForFunctions * @param Realtime $queueForRealtime - * @param StatsUsage $queueForStatsUsage * @param Event $queueForEvents * @param Document $project * @param Document $function @@ -343,9 +337,8 @@ class Functions extends Action Webhook $queueForWebhooks, Func $queueForFunctions, Realtime $queueForRealtime, - StatsUsage $queueForStatsUsage, Event $queueForEvents, - ExecutionEvent $queueForExecutions, + Bus $bus, Document $project, Document $function, Executor $executor, @@ -373,19 +366,19 @@ class Functions extends Action if ($deployment->getAttribute('resourceId') !== $functionId) { $errorMessage = 'The execution could not be completed because a corresponding deployment was not found. A function deployment needs to be created before it can be executed. Please create a deployment for your function and try again.'; - $this->fail($errorMessage, $project, $queueForExecutions, $function, $trigger, $path, $method, $user, $jwt, $event); + $this->fail($errorMessage, $project, $bus, $function, $trigger, $path, $method, $user, $jwt, $event); return; } if ($deployment->isEmpty()) { $errorMessage = 'The execution could not be completed because a corresponding deployment was not found. A function deployment needs to be created before it can be executed. Please create a deployment for your function and try again.'; - $this->fail($errorMessage, $project, $queueForExecutions, $function, $trigger, $path, $method, $user, $jwt, $event); + $this->fail($errorMessage, $project, $bus, $function, $trigger, $path, $method, $user, $jwt, $event); return; } if ($deployment->getAttribute('status') !== 'ready') { $errorMessage = 'The execution could not be completed because the build is not ready. Please wait for the build to complete and try again.'; - $this->fail($errorMessage, $project, $queueForExecutions, $function, $trigger, $path, $method, $user, $jwt, $event); + $this->fail($errorMessage, $project, $bus, $function, $trigger, $path, $method, $user, $jwt, $event); return; } @@ -592,26 +585,12 @@ class Functions extends Action $error = $th->getMessage(); $errorCode = $th->getCode(); } finally { - /** Persist final execution status */ - $queueForExecutions - ->setExecution($execution) - ->setProject($project) - ->trigger(); - - /** Trigger usage queue */ - $queueForStatsUsage - ->setProject($project) - ->addMetric(METRIC_EXECUTIONS, 1) - ->addMetric(str_replace(['{resourceType}'], [RESOURCE_TYPE_FUNCTIONS], METRIC_RESOURCE_TYPE_EXECUTIONS), 1) - ->addMetric(str_replace(['{resourceType}', '{resourceInternalId}'], [RESOURCE_TYPE_FUNCTIONS, $function->getSequence()], METRIC_RESOURCE_TYPE_ID_EXECUTIONS), 1) - ->addMetric(METRIC_EXECUTIONS_COMPUTE, (int)($execution->getAttribute('duration') * 1000))// per project - ->addMetric(str_replace(['{resourceType}'], [RESOURCE_TYPE_FUNCTIONS], METRIC_RESOURCE_TYPE_EXECUTIONS_COMPUTE), (int)($execution->getAttribute('duration') * 1000)) - ->addMetric(str_replace(['{resourceType}', '{resourceInternalId}'], [RESOURCE_TYPE_FUNCTIONS, $function->getSequence()], METRIC_RESOURCE_TYPE_ID_EXECUTIONS_COMPUTE), (int)($execution->getAttribute('duration') * 1000)) - ->addMetric(METRIC_EXECUTIONS_MB_SECONDS, (int)(($spec['memory'] ?? APP_COMPUTE_MEMORY_DEFAULT) * $execution->getAttribute('duration', 0) * ($spec['cpus'] ?? APP_COMPUTE_CPUS_DEFAULT))) - ->addMetric(str_replace(['{resourceType}'], [RESOURCE_TYPE_FUNCTIONS], METRIC_RESOURCE_TYPE_EXECUTIONS_MB_SECONDS), (int)(($spec['memory'] ?? APP_COMPUTE_MEMORY_DEFAULT) * $execution->getAttribute('duration', 0) * ($spec['cpus'] ?? APP_COMPUTE_CPUS_DEFAULT))) - ->addMetric(str_replace(['{resourceType}', '{resourceInternalId}'], [RESOURCE_TYPE_FUNCTIONS, $function->getSequence()], METRIC_RESOURCE_TYPE_ID_EXECUTIONS_MB_SECONDS), (int)(($spec['memory'] ?? APP_COMPUTE_MEMORY_DEFAULT) * $execution->getAttribute('duration', 0) * ($spec['cpus'] ?? APP_COMPUTE_CPUS_DEFAULT))) - ->trigger() - ; + /** Persist final execution status and record usage */ + $bus->dispatch(new ExecutionCompleted( + execution: $execution->getArrayCopy(), + project: $project->getArrayCopy(), + spec: $spec, + )); } $executionModel = new Execution(); diff --git a/src/Utopia/Bus/Bus.php b/src/Utopia/Bus/Bus.php new file mode 100644 index 0000000000..8c4d76958d --- /dev/null +++ b/src/Utopia/Bus/Bus.php @@ -0,0 +1,59 @@ +, Listener[]> */ + private array $listeners = []; + + /** @var ?\Closure(string): mixed */ + private ?\Closure $resolver = null; + + public function setResolver(callable $resolver): self + { + $this->resolver = $resolver(...); + return $this; + } + + public function subscribe(Listener $listener): self + { + foreach ($listener::getEvents() as $event) { + $this->listeners[$event][] = $listener; + } + return $this; + } + + public function dispatch(Event $event): void + { + if ($this->resolver === null) { + throw new \LogicException('Bus resolver must be set via setResolver() before dispatching events'); + } + + $resolver = $this->resolver; + $listeners = $this->listeners[$event::class] ?? []; + + /** @var array}> $resolved */ + $resolved = []; + foreach ($listeners as $listener) { + $deps = array_map($resolver, $listener->getInjections()); + $resolved[] = [$listener, $deps]; + } + + go(function () use ($resolved, $event) { + foreach ($resolved as [$listener, $deps]) { + $action = 'listener.' . $listener::getName(); + Span::init($action); + try { + ($listener->getCallback())($event, ...$deps); + } catch (\Throwable $e) { + Span::error($e); + } finally { + Span::current()?->finish(); + } + } + }); + } +} diff --git a/src/Utopia/Bus/Event.php b/src/Utopia/Bus/Event.php new file mode 100644 index 0000000000..1423f1198d --- /dev/null +++ b/src/Utopia/Bus/Event.php @@ -0,0 +1,7 @@ + */ + protected array $injections = []; + protected ?\Closure $callback = null; + + abstract public static function getName(): string; + + /** + * @return array> + */ + abstract public static function getEvents(): array; + + protected function desc(string $desc): self + { + $this->desc = $desc; + return $this; + } + + protected function inject(string $injection): self + { + $this->injections[] = $injection; + return $this; + } + + protected function callback(callable $callback): self + { + $this->callback = $callback(...); + return $this; + } + + /** @return array */ + public function getInjections(): array + { + return $this->injections; + } + + public function getCallback(): callable + { + if ($this->callback === null) { + throw new \LogicException(static::class . ' must set a callback via $this->callback()'); + } + + return $this->callback; + } +} From 2081c4c42c1bc563b73295289a6b5e2937766dac Mon Sep 17 00:00:00 2001 From: loks0n <22452787+loks0n@users.noreply.github.com> Date: Mon, 23 Feb 2026 20:01:54 +0000 Subject: [PATCH 06/13] refactor: replace bandwidth queueForStatsUsage with Bus events Co-Authored-By: Claude Opus 4.6 --- app/controllers/general.php | 101 ++++++++-------------- app/controllers/shared/api.php | 20 ++--- src/Appwrite/Bus/Listeners/Usage.php | 62 ++++++++++++- src/Appwrite/Bus/RequestCompleted.php | 20 +++++ src/Appwrite/Bus/SiteRequestCompleted.php | 21 +++++ 5 files changed, 144 insertions(+), 80 deletions(-) create mode 100644 src/Appwrite/Bus/RequestCompleted.php create mode 100644 src/Appwrite/Bus/SiteRequestCompleted.php diff --git a/app/controllers/general.php b/app/controllers/general.php index 5f1a5e5d9d..12ab8ea9b4 100644 --- a/app/controllers/general.php +++ b/app/controllers/general.php @@ -9,7 +9,8 @@ use Appwrite\Bus\ExecutionCompleted; use Appwrite\Event\Certificate; use Appwrite\Event\Delete as DeleteEvent; use Appwrite\Event\Event; -use Appwrite\Event\StatsUsage; +use Appwrite\Bus\RequestCompleted; +use Appwrite\Bus\SiteRequestCompleted; use Appwrite\Extend\Exception as AppwriteException; use Appwrite\Network\Cors; use Appwrite\Platform\Appwrite; @@ -63,7 +64,7 @@ Config::setParam('domainVerification', false); Config::setParam('cookieDomain', 'localhost'); Config::setParam('cookieSamesite', Response::COOKIE_SAMESITE_NONE); -function router(Http $utopia, Database $dbForPlatform, callable $getProjectDB, SwooleRequest $swooleRequest, Request $request, Response $response, Log $log, Event $queueForEvents, StatsUsage $queueForStatsUsage, Bus $bus, Executor $executor, Reader $geodb, callable $isResourceBlocked, array $platform, string $previewHostname, Authorization $authorization, ?Key $apiKey, DeleteEvent $queueForDeletes, int $executionsRetentionCount) +function router(Http $utopia, Database $dbForPlatform, callable $getProjectDB, SwooleRequest $swooleRequest, Request $request, Response $response, Log $log, Event $queueForEvents, Bus $bus, Executor $executor, Reader $geodb, callable $isResourceBlocked, array $platform, string $previewHostname, Authorization $authorization, ?Key $apiKey, DeleteEvent $queueForDeletes, int $executionsRetentionCount) { $host = $request->getHostname() ?? ''; if (!empty($previewHostname)) { @@ -757,39 +758,21 @@ function router(Http $utopia, Database $dbForPlatform, callable $getProjectDB, S ->setStatusCode($execution['responseStatusCode'] ?? 200) ->send($body); - $fileSize = 0; - $file = $request->getFiles('file'); - if (!empty($file)) { - $fileSize = (\is_array($file['size']) && isset($file['size'][0])) ? $file['size'][0] : $file['size']; - } - - if (!empty($apiKey) && !empty($apiKey->getDisabledMetrics())) { - foreach ($apiKey->getDisabledMetrics() as $key) { - $queueForStatsUsage->disableMetric($key); - } - } - if ($deployment->getAttribute('resourceType') === 'sites') { - $queueForStatsUsage - ->disableMetric(METRIC_NETWORK_REQUESTS) - ->disableMetric(METRIC_NETWORK_INBOUND) - ->disableMetric(METRIC_NETWORK_OUTBOUND) - ->addMetric(METRIC_SITES_REQUESTS, 1) - ->addMetric(METRIC_SITES_INBOUND, $request->getSize() + $fileSize) - ->addMetric(METRIC_SITES_OUTBOUND, $response->getSize()) - ->addMetric(str_replace('{siteInternalId}', $resource->getSequence(), METRIC_SITES_ID_REQUESTS), 1) - ->addMetric(str_replace('{siteInternalId}', $resource->getSequence(), METRIC_SITES_ID_INBOUND), $request->getSize() + $fileSize) - ->addMetric(str_replace('{siteInternalId}', $resource->getSequence(), METRIC_SITES_ID_OUTBOUND), $response->getSize()) - ; + $bus->dispatch(new SiteRequestCompleted( + project: $project->getArrayCopy(), + request: $request, + response: $response, + siteInternalId: $resource->getSequence(), + )); + } else { + $bus->dispatch(new RequestCompleted( + project: $project->getArrayCopy(), + request: $request, + response: $response, + )); } - $queueForStatsUsage - ->addMetric(METRIC_NETWORK_REQUESTS, 1) - ->addMetric(METRIC_NETWORK_INBOUND, $request->getSize() + $fileSize) - ->addMetric(METRIC_NETWORK_OUTBOUND, $response->getSize()) - ->setProject($project) - ->trigger(); - /* cleanup */ if ($executionsRetentionCount > 0 && ENABLE_EXECUTIONS_LIMIT_ON_ROUTE) { $resourceType = $type === 'function' @@ -852,7 +835,6 @@ Http::init() ->inject('locale') ->inject('localeCodes') ->inject('geodb') - ->inject('queueForStatsUsage') ->inject('queueForEvents') ->inject('bus') ->inject('executor') @@ -865,7 +847,7 @@ Http::init() ->inject('authorization') ->inject('queueForDeletes') ->inject('executionsRetentionCount') - ->action(function (Http $utopia, SwooleRequest $swooleRequest, Request $request, Response $response, Log $log, Document $project, Database $dbForPlatform, callable $getProjectDB, Locale $locale, array $localeCodes, Reader $geodb, StatsUsage $queueForStatsUsage, Event $queueForEvents, Bus $bus, Executor $executor, array $platform, callable $isResourceBlocked, string $previewHostname, Document $devKey, ?Key $apiKey, Cors $cors, Authorization $authorization, DeleteEvent $queueForDeletes, int $executionsRetentionCount) { + ->action(function (Http $utopia, SwooleRequest $swooleRequest, Request $request, Response $response, Log $log, Document $project, Database $dbForPlatform, callable $getProjectDB, Locale $locale, array $localeCodes, Reader $geodb, Event $queueForEvents, Bus $bus, Executor $executor, array $platform, callable $isResourceBlocked, string $previewHostname, Document $devKey, ?Key $apiKey, Cors $cors, Authorization $authorization, DeleteEvent $queueForDeletes, int $executionsRetentionCount) { /* * Appwrite Router */ @@ -873,7 +855,7 @@ Http::init() $platformHostnames = $platform['hostnames'] ?? []; // Only run Router when external domain if (!\in_array($hostname, $platformHostnames) || !empty($previewHostname)) { - if (router($utopia, $dbForPlatform, $getProjectDB, $swooleRequest, $request, $response, $log, $queueForEvents, $queueForStatsUsage, $bus, $executor, $geodb, $isResourceBlocked, $platform, $previewHostname, $authorization, $apiKey, $queueForDeletes, $executionsRetentionCount)) { + if (router($utopia, $dbForPlatform, $getProjectDB, $swooleRequest, $request, $response, $log, $queueForEvents, $bus, $executor, $geodb, $isResourceBlocked, $platform, $previewHostname, $authorization, $apiKey, $queueForDeletes, $executionsRetentionCount)) { $utopia->getRoute()?->label('router', true); } } @@ -1149,7 +1131,6 @@ Http::options() ->inject('dbForPlatform') ->inject('getProjectDB') ->inject('queueForEvents') - ->inject('queueForStatsUsage') ->inject('bus') ->inject('executor') ->inject('geodb') @@ -1163,14 +1144,14 @@ Http::options() ->inject('authorization') ->inject('queueForDeletes') ->inject('executionsRetentionCount') - ->action(function (Http $utopia, SwooleRequest $swooleRequest, Request $request, Response $response, Log $log, Database $dbForPlatform, callable $getProjectDB, Event $queueForEvents, StatsUsage $queueForStatsUsage, Bus $bus, Executor $executor, Reader $geodb, callable $isResourceBlocked, array $platform, string $previewHostname, Document $project, Document $devKey, ?Key $apiKey, Cors $cors, Authorization $authorization, DeleteEvent $queueForDeletes, int $executionsRetentionCount) { + ->action(function (Http $utopia, SwooleRequest $swooleRequest, Request $request, Response $response, Log $log, Database $dbForPlatform, callable $getProjectDB, Event $queueForEvents, Bus $bus, Executor $executor, Reader $geodb, callable $isResourceBlocked, array $platform, string $previewHostname, Document $project, Document $devKey, ?Key $apiKey, Cors $cors, Authorization $authorization, DeleteEvent $queueForDeletes, int $executionsRetentionCount) { /* * Appwrite Router */ $platformHostnames = $platform['hostnames'] ?? []; // Only run Router when external domain if (!in_array($request->getHostname(), $platformHostnames) || !empty($previewHostname)) { - if (router($utopia, $dbForPlatform, $getProjectDB, $swooleRequest, $request, $response, $log, $queueForEvents, $queueForStatsUsage, $bus, $executor, $geodb, $isResourceBlocked, $platform, $previewHostname, $authorization, $apiKey, $queueForDeletes, $executionsRetentionCount)) { + if (router($utopia, $dbForPlatform, $getProjectDB, $swooleRequest, $request, $response, $log, $queueForEvents, $bus, $executor, $geodb, $isResourceBlocked, $platform, $previewHostname, $authorization, $apiKey, $queueForDeletes, $executionsRetentionCount)) { $utopia->getRoute()?->label('router', true); } } @@ -1186,12 +1167,11 @@ Http::options() /** OPTIONS requests in utopia do not execute shutdown handlers, as a result we need to track the OPTIONS requests explicitly * @see https://github.com/utopia-php/http/blob/0.33.16/src/App.php#L825-L855 */ - $queueForStatsUsage - ->addMetric(METRIC_NETWORK_REQUESTS, 1) - ->addMetric(METRIC_NETWORK_INBOUND, $request->getSize()) - ->addMetric(METRIC_NETWORK_OUTBOUND, $response->getSize()) - ->setProject($project) - ->trigger(); + $bus->dispatch(new RequestCompleted( + project: $project->getArrayCopy(), + request: $request, + response: $response, + )); }); Http::error() @@ -1202,10 +1182,10 @@ Http::error() ->inject('project') ->inject('logger') ->inject('log') - ->inject('queueForStatsUsage') + ->inject('bus') ->inject('devKey') ->inject('authorization') - ->action(function (Throwable $error, Http $utopia, Request $request, Response $response, Document $project, ?Logger $logger, Log $log, StatsUsage $queueForStatsUsage, Document $devKey, Authorization $authorization) { + ->action(function (Throwable $error, Http $utopia, Request $request, Response $response, Document $project, ?Logger $logger, Log $log, Bus $bus, Document $devKey, Authorization $authorization) { $version = System::getEnv('_APP_VERSION', 'UNKNOWN'); $route = $utopia->getRoute(); $class = \get_class($error); @@ -1278,21 +1258,12 @@ Http::error() */ if (!$publish && $project->getId() !== 'console') { if (!DBUser::isPrivileged($authorization->getRoles())) { - $fileSize = 0; - $file = $request->getFiles('file'); - if (!empty($file)) { - $fileSize = (\is_array($file['size']) && isset($file['size'][0])) ? $file['size'][0] : $file['size']; - } - - $queueForStatsUsage - ->addMetric(METRIC_NETWORK_REQUESTS, 1) - ->addMetric(METRIC_NETWORK_INBOUND, $request->getSize() + $fileSize) - ->addMetric(METRIC_NETWORK_OUTBOUND, $response->getSize()); + $bus->dispatch(new RequestCompleted( + project: $project->getArrayCopy(), + request: $request, + response: $response, + )); } - - $queueForStatsUsage - ->setProject($project) - ->trigger(); } if ($logger && $publish) { @@ -1539,7 +1510,6 @@ Http::get('/robots.txt') ->inject('dbForPlatform') ->inject('getProjectDB') ->inject('queueForEvents') - ->inject('queueForStatsUsage') ->inject('bus') ->inject('executor') ->inject('geodb') @@ -1550,13 +1520,13 @@ Http::get('/robots.txt') ->inject('authorization') ->inject('queueForDeletes') ->inject('executionsRetentionCount') - ->action(function (Http $utopia, SwooleRequest $swooleRequest, Request $request, Response $response, Log $log, Database $dbForPlatform, callable $getProjectDB, Event $queueForEvents, StatsUsage $queueForStatsUsage, Bus $bus, Executor $executor, Reader $geodb, callable $isResourceBlocked, array $platform, string $previewHostname, ?Key $apiKey, Authorization $authorization, DeleteEvent $queueForDeletes, int $executionsRetentionCount) { + ->action(function (Http $utopia, SwooleRequest $swooleRequest, Request $request, Response $response, Log $log, Database $dbForPlatform, callable $getProjectDB, Event $queueForEvents, Bus $bus, Executor $executor, Reader $geodb, callable $isResourceBlocked, array $platform, string $previewHostname, ?Key $apiKey, Authorization $authorization, DeleteEvent $queueForDeletes, int $executionsRetentionCount) { $platformHostnames = $platform['hostnames'] ?? []; if (in_array($request->getHostname(), $platformHostnames) || !empty($previewHostname)) { $template = new View(__DIR__ . '/../views/general/robots.phtml'); $response->text($template->render(false)); } else { - if (router($utopia, $dbForPlatform, $getProjectDB, $swooleRequest, $request, $response, $log, $queueForEvents, $queueForStatsUsage, $bus, $executor, $geodb, $isResourceBlocked, $platform, $previewHostname, $authorization, $apiKey, $queueForDeletes, $executionsRetentionCount)) { + if (router($utopia, $dbForPlatform, $getProjectDB, $swooleRequest, $request, $response, $log, $queueForEvents, $bus, $executor, $geodb, $isResourceBlocked, $platform, $previewHostname, $authorization, $apiKey, $queueForDeletes, $executionsRetentionCount)) { $utopia->getRoute()?->label('router', true); } } @@ -1574,7 +1544,6 @@ Http::get('/humans.txt') ->inject('dbForPlatform') ->inject('getProjectDB') ->inject('queueForEvents') - ->inject('queueForStatsUsage') ->inject('bus') ->inject('executor') ->inject('geodb') @@ -1585,13 +1554,13 @@ Http::get('/humans.txt') ->inject('authorization') ->inject('queueForDeletes') ->inject('executionsRetentionCount') - ->action(function (Http $utopia, SwooleRequest $swooleRequest, Request $request, Response $response, Log $log, Database $dbForPlatform, callable $getProjectDB, Event $queueForEvents, StatsUsage $queueForStatsUsage, Bus $bus, Executor $executor, Reader $geodb, callable $isResourceBlocked, array $platform, string $previewHostname, ?Key $apiKey, Authorization $authorization, DeleteEvent $queueForDeletes, int $executionsRetentionCount) { + ->action(function (Http $utopia, SwooleRequest $swooleRequest, Request $request, Response $response, Log $log, Database $dbForPlatform, callable $getProjectDB, Event $queueForEvents, Bus $bus, Executor $executor, Reader $geodb, callable $isResourceBlocked, array $platform, string $previewHostname, ?Key $apiKey, Authorization $authorization, DeleteEvent $queueForDeletes, int $executionsRetentionCount) { $platformHostnames = $platform['hostnames'] ?? []; if (in_array($request->getHostname(), $platformHostnames) || !empty($previewHostname)) { $template = new View(__DIR__ . '/../views/general/humans.phtml'); $response->text($template->render(false)); } else { - if (router($utopia, $dbForPlatform, $getProjectDB, $swooleRequest, $request, $response, $log, $queueForEvents, $queueForStatsUsage, $bus, $executor, $geodb, $isResourceBlocked, $platform, $previewHostname, $authorization, $apiKey, $queueForDeletes, $executionsRetentionCount)) { + if (router($utopia, $dbForPlatform, $getProjectDB, $swooleRequest, $request, $response, $log, $queueForEvents, $bus, $executor, $geodb, $isResourceBlocked, $platform, $previewHostname, $authorization, $apiKey, $queueForDeletes, $executionsRetentionCount)) { $utopia->getRoute()?->label('router', true); } } diff --git a/app/controllers/shared/api.php b/app/controllers/shared/api.php index c6499ff9b6..a0ae716559 100644 --- a/app/controllers/shared/api.php +++ b/app/controllers/shared/api.php @@ -2,6 +2,7 @@ use Appwrite\Auth\Key; use Appwrite\Auth\MFA\Type\TOTP; +use Appwrite\Bus\RequestCompleted; use Appwrite\Event\Audit; use Appwrite\Event\Build; use Appwrite\Event\Database as EventDatabase; @@ -21,6 +22,7 @@ use Appwrite\Utopia\Database\Documents\User; use Appwrite\Utopia\Request; use Appwrite\Utopia\Response; use Utopia\Abuse\Abuse; +use Utopia\Bus\Bus; use Utopia\Cache\Adapter\Filesystem; use Utopia\Cache\Cache; use Utopia\Config\Config; @@ -746,7 +748,8 @@ Http::shutdown() ->inject('authorization') ->inject('timelimit') ->inject('eventProcessor') - ->action(function (Http $utopia, Request $request, Response $response, Document $project, User $user, Event $queueForEvents, Audit $queueForAudits, StatsUsage $queueForStatsUsage, Delete $queueForDeletes, EventDatabase $queueForDatabase, Build $queueForBuilds, Messaging $queueForMessaging, Func $queueForFunctions, Event $queueForWebhooks, Realtime $queueForRealtime, Database $dbForProject, Authorization $authorization, callable $timelimit, EventProcessor $eventProcessor) use ($parseLabel) { + ->inject('bus') + ->action(function (Http $utopia, Request $request, Response $response, Document $project, User $user, Event $queueForEvents, Audit $queueForAudits, StatsUsage $queueForStatsUsage, Delete $queueForDeletes, EventDatabase $queueForDatabase, Build $queueForBuilds, Messaging $queueForMessaging, Func $queueForFunctions, Event $queueForWebhooks, Realtime $queueForRealtime, Database $dbForProject, Authorization $authorization, callable $timelimit, EventProcessor $eventProcessor, Bus $bus) use ($parseLabel) { $responsePayload = $response->getPayload(); @@ -958,16 +961,11 @@ Http::shutdown() if ($project->getId() !== 'console') { if (!User::isPrivileged($authorization->getRoles())) { - $fileSize = 0; - $file = $request->getFiles('file'); - if (!empty($file)) { - $fileSize = (\is_array($file['size']) && isset($file['size'][0])) ? $file['size'][0] : $file['size']; - } - - $queueForStatsUsage - ->addMetric(METRIC_NETWORK_REQUESTS, 1) - ->addMetric(METRIC_NETWORK_INBOUND, $request->getSize() + $fileSize) - ->addMetric(METRIC_NETWORK_OUTBOUND, $response->getSize()); + $bus->dispatch(new RequestCompleted( + project: $project->getArrayCopy(), + request: $request, + response: $response, + )); } $queueForStatsUsage diff --git a/src/Appwrite/Bus/Listeners/Usage.php b/src/Appwrite/Bus/Listeners/Usage.php index e91bd910f5..7055c01b7e 100644 --- a/src/Appwrite/Bus/Listeners/Usage.php +++ b/src/Appwrite/Bus/Listeners/Usage.php @@ -3,7 +3,12 @@ namespace Appwrite\Bus\Listeners; use Appwrite\Bus\ExecutionCompleted; +use Appwrite\Bus\RequestCompleted; +use Appwrite\Bus\SiteRequestCompleted; use Appwrite\Event\StatsUsage; +use Appwrite\Utopia\Request; +use Appwrite\Utopia\Response; +use Utopia\Bus\Event; use Utopia\Bus\Listener; use Utopia\Database\Document; use Utopia\Queue\Publisher; @@ -17,18 +22,32 @@ class Usage extends Listener public static function getEvents(): array { - return [ExecutionCompleted::class]; + return [ + ExecutionCompleted::class, + RequestCompleted::class, + SiteRequestCompleted::class, + ]; } public function __construct() { $this - ->desc('Records execution usage metrics') + ->desc('Records usage metrics') ->inject('publisher') ->callback($this->handle(...)); } - public function handle(ExecutionCompleted $event, Publisher $publisher): void + public function handle(Event $event, Publisher $publisher): void + { + match (true) { + $event instanceof ExecutionCompleted => $this->handleExecutionCompleted($event, $publisher), + $event instanceof SiteRequestCompleted => $this->handleSiteRequestCompleted($event, $publisher), + $event instanceof RequestCompleted => $this->handleRequestCompleted($event, $publisher), + default => null, + }; + } + + private function handleExecutionCompleted(ExecutionCompleted $event, Publisher $publisher): void { $execution = new Document($event->execution); $resource = new Document($event->resource); @@ -61,4 +80,41 @@ class Usage extends Listener ->addMetric(str_replace(['{resourceType}', '{resourceInternalId}'], [$resourceType, $resourceInternalId], METRIC_RESOURCE_TYPE_ID_EXECUTIONS_MB_SECONDS), $mbSeconds) ->trigger(); } + + private function handleRequestCompleted(RequestCompleted $event, Publisher $publisher): void + { + $fileSize = 0; + $file = $event->request->getFiles('file'); + if (!empty($file)) { + $fileSize = (\is_array($file['size']) && isset($file['size'][0])) ? $file['size'][0] : $file['size']; + } + + $queueForStatsUsage = new StatsUsage($publisher); + $queueForStatsUsage + ->setProject(new Document($event->project)) + ->addMetric(METRIC_NETWORK_REQUESTS, 1) + ->addMetric(METRIC_NETWORK_INBOUND, $event->request->getSize() + $fileSize) + ->addMetric(METRIC_NETWORK_OUTBOUND, $event->response->getSize()) + ->trigger(); + } + + private function handleSiteRequestCompleted(SiteRequestCompleted $event, Publisher $publisher): void + { + $fileSize = 0; + $file = $event->request->getFiles('file'); + if (!empty($file)) { + $fileSize = (\is_array($file['size']) && isset($file['size'][0])) ? $file['size'][0] : $file['size']; + } + + $queueForStatsUsage = new StatsUsage($publisher); + $queueForStatsUsage + ->setProject(new Document($event->project)) + ->addMetric(METRIC_SITES_REQUESTS, 1) + ->addMetric(METRIC_SITES_INBOUND, $event->request->getSize() + $fileSize) + ->addMetric(METRIC_SITES_OUTBOUND, $event->response->getSize()) + ->addMetric(str_replace('{siteInternalId}', $event->siteInternalId, METRIC_SITES_ID_REQUESTS), 1) + ->addMetric(str_replace('{siteInternalId}', $event->siteInternalId, METRIC_SITES_ID_INBOUND), $event->request->getSize() + $fileSize) + ->addMetric(str_replace('{siteInternalId}', $event->siteInternalId, METRIC_SITES_ID_OUTBOUND), $event->response->getSize()) + ->trigger(); + } } diff --git a/src/Appwrite/Bus/RequestCompleted.php b/src/Appwrite/Bus/RequestCompleted.php new file mode 100644 index 0000000000..773d93e873 --- /dev/null +++ b/src/Appwrite/Bus/RequestCompleted.php @@ -0,0 +1,20 @@ + $project + */ + public function __construct( + public readonly array $project, + public readonly Request $request, + public readonly Response $response, + ) { + } +} diff --git a/src/Appwrite/Bus/SiteRequestCompleted.php b/src/Appwrite/Bus/SiteRequestCompleted.php new file mode 100644 index 0000000000..86fd220387 --- /dev/null +++ b/src/Appwrite/Bus/SiteRequestCompleted.php @@ -0,0 +1,21 @@ + $project + */ + public function __construct( + public readonly array $project, + public readonly Request $request, + public readonly Response $response, + public readonly string $siteInternalId = '', + ) { + } +} From c171e0c3a21f7a8fb14f69f40ca0130ab9cf6212 Mon Sep 17 00:00:00 2001 From: loks0n <22452787+loks0n@users.noreply.github.com> Date: Mon, 23 Feb 2026 20:10:57 +0000 Subject: [PATCH 07/13] refactor: add bus.event span attribute to listener invocations Co-Authored-By: Claude Opus 4.6 --- src/Utopia/Bus/Bus.php | 1 + 1 file changed, 1 insertion(+) diff --git a/src/Utopia/Bus/Bus.php b/src/Utopia/Bus/Bus.php index 8c4d76958d..f8cba49a97 100644 --- a/src/Utopia/Bus/Bus.php +++ b/src/Utopia/Bus/Bus.php @@ -46,6 +46,7 @@ class Bus foreach ($resolved as [$listener, $deps]) { $action = 'listener.' . $listener::getName(); Span::init($action); + Span::add('bus.event', $event::class); try { ($listener->getCallback())($event, ...$deps); } catch (\Throwable $e) { From 20f248a6ae17c8685b5a8e3f2099b3e6f186ea29 Mon Sep 17 00:00:00 2001 From: loks0n <22452787+loks0n@users.noreply.github.com> Date: Mon, 23 Feb 2026 20:35:38 +0000 Subject: [PATCH 08/13] refactor: consolidate SiteRequestCompleted into RequestCompleted with optional deployment Co-Authored-By: Claude Opus 4.6 --- app/controllers/general.php | 25 +++------ app/controllers/shared/api.php | 2 +- .../Bus/{ => Events}/ExecutionCompleted.php | 2 +- .../Bus/{ => Events}/RequestCompleted.php | 4 +- src/Appwrite/Bus/Listeners/Log.php | 2 +- src/Appwrite/Bus/Listeners/Usage.php | 54 +++++++++---------- src/Appwrite/Bus/SiteRequestCompleted.php | 21 -------- src/Appwrite/Platform/Workers/Functions.php | 2 +- 8 files changed, 39 insertions(+), 73 deletions(-) rename src/Appwrite/Bus/{ => Events}/ExecutionCompleted.php (94%) rename src/Appwrite/Bus/{ => Events}/RequestCompleted.php (74%) delete mode 100644 src/Appwrite/Bus/SiteRequestCompleted.php diff --git a/app/controllers/general.php b/app/controllers/general.php index 12ab8ea9b4..15f16de9c5 100644 --- a/app/controllers/general.php +++ b/app/controllers/general.php @@ -5,12 +5,11 @@ require_once __DIR__ . '/../init.php'; use Ahc\Jwt\JWT; use Ahc\Jwt\JWTException; use Appwrite\Auth\Key; -use Appwrite\Bus\ExecutionCompleted; +use Appwrite\Bus\Events\ExecutionCompleted; use Appwrite\Event\Certificate; use Appwrite\Event\Delete as DeleteEvent; use Appwrite\Event\Event; -use Appwrite\Bus\RequestCompleted; -use Appwrite\Bus\SiteRequestCompleted; +use Appwrite\Bus\Events\RequestCompleted; use Appwrite\Extend\Exception as AppwriteException; use Appwrite\Network\Cors; use Appwrite\Platform\Appwrite; @@ -758,20 +757,12 @@ function router(Http $utopia, Database $dbForPlatform, callable $getProjectDB, S ->setStatusCode($execution['responseStatusCode'] ?? 200) ->send($body); - if ($deployment->getAttribute('resourceType') === 'sites') { - $bus->dispatch(new SiteRequestCompleted( - project: $project->getArrayCopy(), - request: $request, - response: $response, - siteInternalId: $resource->getSequence(), - )); - } else { - $bus->dispatch(new RequestCompleted( - project: $project->getArrayCopy(), - request: $request, - response: $response, - )); - } + $bus->dispatch(new RequestCompleted( + project: $project->getArrayCopy(), + request: $request, + response: $response, + deployment: $deployment->getArrayCopy(), + )); /* cleanup */ if ($executionsRetentionCount > 0 && ENABLE_EXECUTIONS_LIMIT_ON_ROUTE) { diff --git a/app/controllers/shared/api.php b/app/controllers/shared/api.php index a0ae716559..d7b2c7339e 100644 --- a/app/controllers/shared/api.php +++ b/app/controllers/shared/api.php @@ -2,7 +2,7 @@ use Appwrite\Auth\Key; use Appwrite\Auth\MFA\Type\TOTP; -use Appwrite\Bus\RequestCompleted; +use Appwrite\Bus\Events\RequestCompleted; use Appwrite\Event\Audit; use Appwrite\Event\Build; use Appwrite\Event\Database as EventDatabase; diff --git a/src/Appwrite/Bus/ExecutionCompleted.php b/src/Appwrite/Bus/Events/ExecutionCompleted.php similarity index 94% rename from src/Appwrite/Bus/ExecutionCompleted.php rename to src/Appwrite/Bus/Events/ExecutionCompleted.php index 1508266ea0..58c84c82f0 100644 --- a/src/Appwrite/Bus/ExecutionCompleted.php +++ b/src/Appwrite/Bus/Events/ExecutionCompleted.php @@ -1,6 +1,6 @@ $project + * @param array $deployment */ public function __construct( public readonly array $project, public readonly Request $request, public readonly Response $response, + public readonly array $deployment = [], ) { } } diff --git a/src/Appwrite/Bus/Listeners/Log.php b/src/Appwrite/Bus/Listeners/Log.php index 12ca7303c8..9bd539d5fe 100644 --- a/src/Appwrite/Bus/Listeners/Log.php +++ b/src/Appwrite/Bus/Listeners/Log.php @@ -2,7 +2,7 @@ namespace Appwrite\Bus\Listeners; -use Appwrite\Bus\ExecutionCompleted; +use Appwrite\Bus\Events\ExecutionCompleted; use Appwrite\Event\Execution; use Utopia\Bus\Listener; use Utopia\Database\Document; diff --git a/src/Appwrite/Bus/Listeners/Usage.php b/src/Appwrite/Bus/Listeners/Usage.php index 7055c01b7e..ef288e427c 100644 --- a/src/Appwrite/Bus/Listeners/Usage.php +++ b/src/Appwrite/Bus/Listeners/Usage.php @@ -2,12 +2,9 @@ namespace Appwrite\Bus\Listeners; -use Appwrite\Bus\ExecutionCompleted; -use Appwrite\Bus\RequestCompleted; -use Appwrite\Bus\SiteRequestCompleted; +use Appwrite\Bus\Events\ExecutionCompleted; +use Appwrite\Bus\Events\RequestCompleted; use Appwrite\Event\StatsUsage; -use Appwrite\Utopia\Request; -use Appwrite\Utopia\Response; use Utopia\Bus\Event; use Utopia\Bus\Listener; use Utopia\Database\Document; @@ -25,7 +22,6 @@ class Usage extends Listener return [ ExecutionCompleted::class, RequestCompleted::class, - SiteRequestCompleted::class, ]; } @@ -41,7 +37,6 @@ class Usage extends Listener { match (true) { $event instanceof ExecutionCompleted => $this->handleExecutionCompleted($event, $publisher), - $event instanceof SiteRequestCompleted => $this->handleSiteRequestCompleted($event, $publisher), $event instanceof RequestCompleted => $this->handleRequestCompleted($event, $publisher), default => null, }; @@ -89,32 +84,31 @@ class Usage extends Listener $fileSize = (\is_array($file['size']) && isset($file['size'][0])) ? $file['size'][0] : $file['size']; } + $project = new Document($event->project); + $deployment = new Document($event->deployment); $queueForStatsUsage = new StatsUsage($publisher); - $queueForStatsUsage - ->setProject(new Document($event->project)) - ->addMetric(METRIC_NETWORK_REQUESTS, 1) - ->addMetric(METRIC_NETWORK_INBOUND, $event->request->getSize() + $fileSize) - ->addMetric(METRIC_NETWORK_OUTBOUND, $event->response->getSize()) - ->trigger(); - } - private function handleSiteRequestCompleted(SiteRequestCompleted $event, Publisher $publisher): void - { - $fileSize = 0; - $file = $event->request->getFiles('file'); - if (!empty($file)) { - $fileSize = (\is_array($file['size']) && isset($file['size'][0])) ? $file['size'][0] : $file['size']; + $inbound = $event->request->getSize() + $fileSize; + $outbound = $event->response->getSize(); + + $queueForStatsUsage->setProject($project); + + if ($deployment->getAttribute('resourceType') === 'sites') { + $siteInternalId = $deployment->getAttribute('resourceInternalId', ''); + $queueForStatsUsage + ->addMetric(METRIC_SITES_REQUESTS, 1) + ->addMetric(METRIC_SITES_INBOUND, $inbound) + ->addMetric(METRIC_SITES_OUTBOUND, $outbound) + ->addMetric(str_replace('{siteInternalId}', $siteInternalId, METRIC_SITES_ID_REQUESTS), 1) + ->addMetric(str_replace('{siteInternalId}', $siteInternalId, METRIC_SITES_ID_INBOUND), $inbound) + ->addMetric(str_replace('{siteInternalId}', $siteInternalId, METRIC_SITES_ID_OUTBOUND), $outbound); + } else { + $queueForStatsUsage + ->addMetric(METRIC_NETWORK_REQUESTS, 1) + ->addMetric(METRIC_NETWORK_INBOUND, $inbound) + ->addMetric(METRIC_NETWORK_OUTBOUND, $outbound); } - $queueForStatsUsage = new StatsUsage($publisher); - $queueForStatsUsage - ->setProject(new Document($event->project)) - ->addMetric(METRIC_SITES_REQUESTS, 1) - ->addMetric(METRIC_SITES_INBOUND, $event->request->getSize() + $fileSize) - ->addMetric(METRIC_SITES_OUTBOUND, $event->response->getSize()) - ->addMetric(str_replace('{siteInternalId}', $event->siteInternalId, METRIC_SITES_ID_REQUESTS), 1) - ->addMetric(str_replace('{siteInternalId}', $event->siteInternalId, METRIC_SITES_ID_INBOUND), $event->request->getSize() + $fileSize) - ->addMetric(str_replace('{siteInternalId}', $event->siteInternalId, METRIC_SITES_ID_OUTBOUND), $event->response->getSize()) - ->trigger(); + $queueForStatsUsage->trigger(); } } diff --git a/src/Appwrite/Bus/SiteRequestCompleted.php b/src/Appwrite/Bus/SiteRequestCompleted.php deleted file mode 100644 index 86fd220387..0000000000 --- a/src/Appwrite/Bus/SiteRequestCompleted.php +++ /dev/null @@ -1,21 +0,0 @@ - $project - */ - public function __construct( - public readonly array $project, - public readonly Request $request, - public readonly Response $response, - public readonly string $siteInternalId = '', - ) { - } -} diff --git a/src/Appwrite/Platform/Workers/Functions.php b/src/Appwrite/Platform/Workers/Functions.php index 6c59bc8da0..18ab087966 100644 --- a/src/Appwrite/Platform/Workers/Functions.php +++ b/src/Appwrite/Platform/Workers/Functions.php @@ -3,7 +3,7 @@ namespace Appwrite\Platform\Workers; use Ahc\Jwt\JWT; -use Appwrite\Bus\ExecutionCompleted; +use Appwrite\Bus\Events\ExecutionCompleted; use Appwrite\Event\Event; use Appwrite\Event\Func; use Appwrite\Event\Realtime; From a0854e05919af810e1538c7e1be78fda35d9beeb Mon Sep 17 00:00:00 2001 From: loks0n <22452787+loks0n@users.noreply.github.com> Date: Tue, 3 Mar 2026 20:06:06 +0000 Subject: [PATCH 09/13] refactor: make Bus dispatch synchronous Remove async coroutine wrapper from event dispatch to simplify execution model and improve trace hierarchy. Listeners now execute synchronously in the caller's context, with dependency resolution inlined. Co-Authored-By: Claude Sonnet 4.5 --- src/Utopia/Bus/Bus.php | 27 +++++++++------------------ 1 file changed, 9 insertions(+), 18 deletions(-) diff --git a/src/Utopia/Bus/Bus.php b/src/Utopia/Bus/Bus.php index f8cba49a97..bef39f0481 100644 --- a/src/Utopia/Bus/Bus.php +++ b/src/Utopia/Bus/Bus.php @@ -35,26 +35,17 @@ class Bus $resolver = $this->resolver; $listeners = $this->listeners[$event::class] ?? []; - /** @var array}> $resolved */ - $resolved = []; foreach ($listeners as $listener) { $deps = array_map($resolver, $listener->getInjections()); - $resolved[] = [$listener, $deps]; - } - - go(function () use ($resolved, $event) { - foreach ($resolved as [$listener, $deps]) { - $action = 'listener.' . $listener::getName(); - Span::init($action); - Span::add('bus.event', $event::class); - try { - ($listener->getCallback())($event, ...$deps); - } catch (\Throwable $e) { - Span::error($e); - } finally { - Span::current()?->finish(); - } + Span::init('listener.' . $listener::getName()); + Span::add('bus.event', $event::class); + try { + ($listener->getCallback())($event, ...$deps); + } catch (\Throwable $e) { + Span::error($e); + } finally { + Span::current()?->finish(); } - }); + } } } From c0737439890e4d1290904497b4cc1bc9fd61c6b5 Mon Sep 17 00:00:00 2001 From: loks0n <22452787+loks0n@users.noreply.github.com> Date: Tue, 3 Mar 2026 20:11:34 +0000 Subject: [PATCH 10/13] fix: lint - order imports in general controller Co-Authored-By: Claude Sonnet 4.5 --- app/controllers/general.php | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/app/controllers/general.php b/app/controllers/general.php index 15f16de9c5..43c7e47ca6 100644 --- a/app/controllers/general.php +++ b/app/controllers/general.php @@ -6,10 +6,10 @@ use Ahc\Jwt\JWT; use Ahc\Jwt\JWTException; use Appwrite\Auth\Key; use Appwrite\Bus\Events\ExecutionCompleted; +use Appwrite\Bus\Events\RequestCompleted; use Appwrite\Event\Certificate; use Appwrite\Event\Delete as DeleteEvent; use Appwrite\Event\Event; -use Appwrite\Bus\Events\RequestCompleted; use Appwrite\Extend\Exception as AppwriteException; use Appwrite\Network\Cors; use Appwrite\Platform\Appwrite; From cf7710b58082171730feb44c393648e3c9fb4b4e Mon Sep 17 00:00:00 2001 From: loks0n <22452787+loks0n@users.noreply.github.com> Date: Tue, 3 Mar 2026 22:47:28 +0000 Subject: [PATCH 11/13] perf: reduce WebSocket timeout in realtime tests from 45s to 2s Realtime E2E tests were taking 24+ minutes due to intentional timeout waits. Many tests verify event filtering by expecting a TimeoutException, and each was waiting 45 seconds. Changes: - Reduce default WebSocket timeout from 45s to 2s in RealtimeBase - Add optional timeout parameter to getWebsocket() methods - Use longer timeouts (5-10s) for tests that legitimately wait for slow operations (function executions, test channel events) - Use named parameters for improved readability Performance impact: - Before: 24:07 minutes (1,447 seconds) - After: 1:09 minutes (69 seconds) - Speedup: ~21x faster Co-Authored-By: Claude Sonnet 4.5 --- tests/e2e/Services/Realtime/RealtimeBase.php | 10 +++--- .../RealtimeCustomClientQueryTest.php | 32 +++++++++++-------- .../Realtime/RealtimeCustomClientTest.php | 12 ++++--- 3 files changed, 33 insertions(+), 21 deletions(-) diff --git a/tests/e2e/Services/Realtime/RealtimeBase.php b/tests/e2e/Services/Realtime/RealtimeBase.php index ee7946d9f7..95f3665e4c 100644 --- a/tests/e2e/Services/Realtime/RealtimeBase.php +++ b/tests/e2e/Services/Realtime/RealtimeBase.php @@ -11,7 +11,8 @@ trait RealtimeBase array $channels = [], array $headers = [], ?string $projectId = null, - ?array $queries = null + ?array $queries = null, + int $timeout = 2 ): WebSocketClient { if (is_null($projectId)) { $projectId = $this->getProject()['$id']; @@ -63,7 +64,7 @@ trait RealtimeBase "ws://appwrite.test/v1/realtime?" . $queryString, [ "headers" => $headers, - "timeout" => 45, + "timeout" => $timeout, ] ); } @@ -74,9 +75,10 @@ trait RealtimeBase * * @param array $queryParams Custom query parameters (e.g., ['channels' => ['project'], 'project' => [...]]) * @param array $headers HTTP headers + * @param int $timeout Timeout in seconds (default: 2) * @return WebSocketClient */ - private function getWebsocketWithCustomQuery(array $queryParams, array $headers = []): WebSocketClient + private function getWebsocketWithCustomQuery(array $queryParams, array $headers = [], int $timeout = 2): WebSocketClient { $queryString = http_build_query($queryParams); @@ -84,7 +86,7 @@ trait RealtimeBase "ws://appwrite.test/v1/realtime?" . $queryString, [ "headers" => $headers, - "timeout" => 45, + "timeout" => $timeout, ] ); } diff --git a/tests/e2e/Services/Realtime/RealtimeCustomClientQueryTest.php b/tests/e2e/Services/Realtime/RealtimeCustomClientQueryTest.php index 8104fa7bd0..30cc70e981 100644 --- a/tests/e2e/Services/Realtime/RealtimeCustomClientQueryTest.php +++ b/tests/e2e/Services/Realtime/RealtimeCustomClientQueryTest.php @@ -2535,29 +2535,35 @@ class RealtimeCustomClientQueryTest extends Scope $projectId = 'console'; // Subscribe without queries - should receive all events - $clientNoQuery = $this->getWebsocket(['tests'], [ - 'origin' => 'http://localhost', - ], $projectId); + $clientNoQuery = $this->getWebsocket( + channels: ['tests'], + headers: ['origin' => 'http://localhost'], + projectId: $projectId, + timeout: 5 + ); $response = json_decode($clientNoQuery->receive(), true); $this->assertEquals('connected', $response['type']); // Subscribe with matching query - should receive events - $clientWithMatchingQuery = $this->getWebsocket(['tests'], [ - 'origin' => 'http://localhost', - ], $projectId, [ - Query::equal('response', ['WS:/v1/realtime:passed'])->toString(), - ]); + $clientWithMatchingQuery = $this->getWebsocket( + channels: ['tests'], + headers: ['origin' => 'http://localhost'], + projectId: $projectId, + queries: [Query::equal('response', ['WS:/v1/realtime:passed'])->toString()], + timeout: 5 + ); $response = json_decode($clientWithMatchingQuery->receive(), true); $this->assertEquals('connected', $response['type']); // Subscribe with non-matching query - should NOT receive events - $clientWithNonMatchingQuery = $this->getWebsocket(['tests'], [ - 'origin' => 'http://localhost', - ], $projectId, [ - Query::equal('response', ['failed'])->toString(), - ]); + $clientWithNonMatchingQuery = $this->getWebsocket( + channels: ['tests'], + headers: ['origin' => 'http://localhost'], + projectId: $projectId, + queries: [Query::equal('response', ['failed'])->toString()] + ); $response = json_decode($clientWithNonMatchingQuery->receive(), true); $this->assertEquals('connected', $response['type']); diff --git a/tests/e2e/Services/Realtime/RealtimeCustomClientTest.php b/tests/e2e/Services/Realtime/RealtimeCustomClientTest.php index b94e382fde..c5fb71fcb6 100644 --- a/tests/e2e/Services/Realtime/RealtimeCustomClientTest.php +++ b/tests/e2e/Services/Realtime/RealtimeCustomClientTest.php @@ -2225,10 +2225,14 @@ class RealtimeCustomClientTest extends Scope $session = $user['session'] ?? ''; $projectId = $this->getProject()['$id']; - $client = $this->getWebsocket(['executions'], [ - 'origin' => 'http://localhost', - 'cookie' => 'a_session_' . $projectId . '=' . $session - ]); + $client = $this->getWebsocket( + channels: ['executions'], + headers: [ + 'origin' => 'http://localhost', + 'cookie' => 'a_session_' . $projectId . '=' . $session + ], + timeout: 10 + ); $response = json_decode($client->receive(), true); From f5047afec9030a49f61d736475395b6f601e8d4d Mon Sep 17 00:00:00 2001 From: fogelito Date: Wed, 4 Mar 2026 10:28:56 +0200 Subject: [PATCH 12/13] Remove index --- app/config/collections/projects.php | 7 ------- 1 file changed, 7 deletions(-) diff --git a/app/config/collections/projects.php b/app/config/collections/projects.php index da37334b4e..984d7ec6b6 100644 --- a/app/config/collections/projects.php +++ b/app/config/collections/projects.php @@ -2190,13 +2190,6 @@ return [ 'lengths' => [], 'orders' => [Database::ORDER_ASC], ], - [ - '$id' => ID::custom('_key_function_internal_id'), // Index not in use remove in the future - 'type' => Database::INDEX_KEY, - 'attributes' => ['functionInternalId'], - 'lengths' => [], - 'orders' => [], - ], [ '$id' => ID::custom('_key_resourceType'), 'type' => Database::INDEX_KEY, From c81029d8aabd929ce4fd066a9176a53195949ea7 Mon Sep 17 00:00:00 2001 From: loks0n <22452787+loks0n@users.noreply.github.com> Date: Wed, 4 Mar 2026 11:17:50 +0000 Subject: [PATCH 13/13] fix: stats publisher --- src/Appwrite/Bus/Listeners/Usage.php | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/Appwrite/Bus/Listeners/Usage.php b/src/Appwrite/Bus/Listeners/Usage.php index ef288e427c..219287033d 100644 --- a/src/Appwrite/Bus/Listeners/Usage.php +++ b/src/Appwrite/Bus/Listeners/Usage.php @@ -29,7 +29,7 @@ class Usage extends Listener { $this ->desc('Records usage metrics') - ->inject('publisher') + ->inject('publisherStatsUsage') ->callback($this->handle(...)); }