WIP: Function proxy

This commit is contained in:
Matej Baco
2022-07-12 13:57:08 +02:00
parent 3d5d45d676
commit 1529c649ff
6 changed files with 194 additions and 79 deletions
+2 -2
View File
@@ -83,5 +83,5 @@ _APP_LOGGING_CONFIG=
DOCKERHUB_PULL_USERNAME=
DOCKERHUB_PULL_PASSWORD=
DOCKERHUB_PULL_EMAIL=
#_APP_EXECUTORS=exec_fra_01=http://appwrite-executor/v1,exec_fra_02=http://appwrite-executor2/v1
_APP_EXECUTORS=exec_fra_01=http://appwrite-executor/v1
_APP_EXECUTORS=executor_01=appwrite-executor,executor_02=appwrite-executor2
#_APP_EXECUTORS=exec_fra_01=http://appwrite-executor/v1
+49 -77
View File
@@ -2,12 +2,12 @@
require_once __DIR__ . '/../vendor/autoload.php';
use FunctionsProxy\Adapter\Random;
use FunctionsProxy\Adapter\RoundRobin;
use Swoole\Coroutine\Http\Client;
use Utopia\Logger\Log;
use Utopia\Logger\Logger;
use Swoole\Http\Server;
use Utopia\Swoole\Request;
use Utopia\Swoole\Response;
use Swoole\Database\RedisConfig;
use Swoole\Database\RedisPool;
use Swoole\Http\Request as SwooleRequest;
@@ -38,23 +38,27 @@ $redisPool = new RedisPool(
64
);
function markOffline(Cache $cache, string $executorId, string $error) {
$adapter = new RoundRobin($redisPool);
function markOffline(Cache $cache, string $executorId, string $error)
{
$data = $cache->load('executors-' . $executorId, 60 * 60 * 24 * 30 * 3); // 3 months
$cache->save('executors-' . $executorId, [ 'state' => 'offline' ]);
$cache->save('executors-' . $executorId, ['status' => 'offline']);
if(!$data || $data['state'] === 'online') {
if (!$data || $data['status'] === 'online') {
Console::warning('Executor "' . $executorId . '" went down! Message:');
Console::warning($error);
}
}
function markOnline(cache $cache, string $executorId) {
function markOnline(cache $cache, string $executorId)
{
$data = $cache->load('executors-' . $executorId, 60 * 60 * 24 * 30 * 3); // 3 months
$cache->save('executors-' . $executorId, [ 'state' => 'online' ]);
$cache->save('executors-' . $executorId, ['status' => 'online']);
if(!$data || $data['state'] === 'offline') {
if (!$data || $data['status'] === 'offline') {
Console::success('Executor "' . $executorId . '" went online.');
}
}
@@ -72,7 +76,7 @@ function fetchExecutorsState(RedisPool $redisPool)
try {
[$id, $hostname] = \explode('=', $executor);
$endpoint = $hostname . '/health';
$endpoint = 'http://' . $hostname . '/v1/health';
$ch = \curl_init();
@@ -166,85 +170,51 @@ Console::success("Waiting for executors to start...");
\sleep(5); // Wait a little so executors can start
// fetchExecutorsState($redisPool);
// Timer::tick(30000, fn (int $timerId, array $params) => fetchExecutorsState($params[0]), [$redisPool]);
Console::success("Functions proxy is ready.");
App::setMode(App::MODE_TYPE_PRODUCTION); // Define Mode
fetchExecutorsState($redisPool);
$http = new Server("0.0.0.0", 80);
/** Set callbacks */
App::error(function ($utopia, $error, $request, $response) {
$route = $utopia->match($request);
logError($error, "httpError", $route);
$run = function (SwooleRequest $request, SwooleResponse $response) use ($adapter) {
$secretKey = $request->header['x-appwrite-functions-proxy-key'] ?? '';
switch ($error->getCode()) {
case 400: // Error allowed publicly
case 401: // Error allowed publicly
case 402: // Error allowed publicly
case 403: // Error allowed publicly
case 404: // Error allowed publicly
case 406: // Error allowed publicly
case 409: // Error allowed publicly
case 412: // Error allowed publicly
case 425: // Error allowed publicly
case 429: // Error allowed publicly
case 501: // Error allowed publicly
case 503: // Error allowed publicly
$code = $error->getCode();
break;
default:
$code = 500; // All other errors get the generic 500 server error status code
if (empty($secretKey)) {
throw new Exception('Missing proxy key');
}
if ($secretKey !== App::getEnv('_APP_FUNCTIONS_PROXY_SECRET', '')) {
throw new Exception('Missing proxy key');
}
$output = [
'message' => $error->getMessage(),
'code' => $error->getCode(),
'file' => $error->getFile(),
'line' => $error->getLine(),
'trace' => $error->getTrace(),
'version' => App::getEnv('_APP_VERSION', 'UNKNOWN')
];
$executorHostname = $adapter->getNextExecutor();
$response
->addHeader('Cache-Control', 'no-cache, no-store, must-revalidate')
->addHeader('Expires', '0')
->addHeader('Pragma', 'no-cache')
->setStatusCode($code);
\var_dump($request->server['request_uri']);
\var_dump($request->server['request_method']);
\var_dump($executorHostname);
$response->json($output);
}, ['utopia', 'error', 'request', 'response']);
$client = new Client($executorHostname, 80);
$client->setMethod($request->server['request_method'] ?? 'GET');
$client->setData($request->getData());
$client->setHeaders(\array_merge($request->header, [
'x-appwrite-executor-key' => App::getEnv('_APP_EXECUTOR_SECRET', '')
]));
App::get('/')
->desc("Proxy request into executor")
->inject('response')
->action(function (Response $response) {
$response
->setStatusCode(Response::STATUS_CODE_OK)
->json([ 'works' => true ]);
});
$status = $client->execute($request->server['request_uri'] ?? '/');
App::init(function ($request, $response) {
$secretKey = $request->getHeader('x-appwrite-functions-proxy-key', '');
$response->setStatusCode(200);
$response->end(\json_encode([
'data' => $status,
'data2' => $client->getBody(),
'data3' => $client->getStatusCode(),
'data4' => $client->errCode
]));
};
if (empty($secretKey)) {
throw new Exception('Missing proxy key', 401);
}
if ($secretKey !== App::getEnv('_APP_FUNCTIONS_PROXY_SECRET', '')) {
throw new Exception('Missing proxy key', 401);
}
}, ['request', 'response']);
$http->on('request', function (SwooleRequest $swooleRequest, SwooleResponse $swooleResponse) {
$request = new Request($swooleRequest);
$response = new Response($swooleResponse);
$app = new App('UTC');
$http->on('start', function () use ($redisPool) {
Timer::tick(30000, fn (int $timerId, array $params) => fetchExecutorsState($params[0]), [$redisPool]);
});
$http->on('request', function (SwooleRequest $swooleRequest, SwooleResponse $swooleResponse) use ($run) {
try {
$app->run($request, $response);
call_user_func($run, $swooleRequest, $swooleResponse);
} catch (\Throwable $th) {
logError($th, "serverError");
$swooleResponse->setStatusCode(500);
@@ -259,4 +229,6 @@ $http->on('request', function (SwooleRequest $swooleRequest, SwooleResponse $swo
}
});
$http->start();
Console::success("Functions proxy is ready.");
$http->start();
+45
View File
@@ -521,6 +521,51 @@ services:
- DOCKERHUB_PULL_USERNAME
- DOCKERHUB_PULL_PASSWORD
appwrite-executor2:
container_name: appwrite-executor2
<<: *x-logging
entrypoint: executor
stop_signal: SIGINT
build:
context: .
args:
- DEBUG=false
- TESTING=true
- VERSION=dev
networks:
appwrite:
runtimes:
volumes:
- /var/run/docker.sock:/var/run/docker.sock
- ./app:/usr/src/code/app
- ./src:/usr/src/code/src
- appwrite-functions:/storage/functions:rw
- appwrite-builds:/storage/builds:rw
- /tmp:/tmp:rw
depends_on:
- redis
- mariadb
- appwrite
- appwrite-functions-proxy
environment:
- _APP_ENV
- _APP_VERSION
- _APP_FUNCTIONS_TIMEOUT
- _APP_FUNCTIONS_BUILD_TIMEOUT
- _APP_FUNCTIONS_CONTAINERS
- _APP_FUNCTIONS_RUNTIMES
- _APP_FUNCTIONS_CPUS
- _APP_FUNCTIONS_MEMORY
- _APP_FUNCTIONS_MEMORY_SWAP
- _APP_FUNCTIONS_INACTIVE_THRESHOLD
- _APP_EXECUTOR_SECRET
- OPEN_RUNTIMES_NETWORK
- _APP_LOGGING_PROVIDER
- _APP_LOGGING_CONFIG
- *x-env-storage
- DOCKERHUB_PULL_USERNAME
- DOCKERHUB_PULL_PASSWORD
appwrite-worker-mails:
entrypoint: worker-mails
<<: *x-logging
+57
View File
@@ -0,0 +1,57 @@
<?php
namespace FunctionsProxy;
use Swoole\Database\RedisPool;
use Utopia\App;
use Utopia\Cache\Adapter\Redis;
use Utopia\Cache\Cache;
abstract class Adapter
{
private RedisPool $redisPool;
public function __construct(RedisPool $redisPool)
{
$this->redisPool = $redisPool;
}
private function getConnection(): array {
$redis = $this->redisPool->get();
$cache = new Cache(new Redis($redis));
return [$cache, fn() => $this->redisPool->put($redis)];
}
protected function getExecutors(): array {
[$cache, $returnCache] = $this->getConnection();
$responseExecutors = [];
try {
$executors = \explode(',', App::getEnv('_APP_EXECUTORS', ''));
foreach ($executors as $executor) {
[$id, $hostname] = \explode('=', $executor);
$data = $cache->load('executors-' . $id, 60 * 60 * 24 * 30 * 3); // 3 months
if($data['status'] !== 'online') {
continue;
}
$responseExecutors[] = [
'id' => $id,
'hostname' => $hostname,
'state' => $data
];
}
} finally {
call_user_func($returnCache);
}
return $responseExecutors;
}
abstract public function getNextExecutor(): string;
}
+15
View File
@@ -0,0 +1,15 @@
<?php
namespace FunctionsProxy\Adapter;
use FunctionsProxy\Adapter;
class Random extends Adapter
{
public function getNextExecutor(): string {
$executors = $this->getExecutors();
$executor = $executors[\array_rand($executors)] ?? null;
return $executor['hostname'] ?? null;
}
}
+26
View File
@@ -0,0 +1,26 @@
<?php
namespace FunctionsProxy\Adapter;
use FunctionsProxy\Adapter;
class RoundRobin extends Adapter
{
private $currentIndex = 0; // TODO: Put into redis to share across proxies
public function getNextExecutor(): string
{
$executors = $this->getExecutors();
$executor = $executors[$this->currentIndex] ?? null;
$this->currentIndex++;
if (!$executor) {
$this->currentIndex = 0;
$executor = $executors[$this->currentIndex] ?? null;
$this->currentIndex++;
}
return $executor['hostname'] ?? null;
}
}