diff --git a/ProcessMaker/Enums/ScriptExecutorType.php b/ProcessMaker/Enums/ScriptExecutorType.php index c4e7774c90..9a217ea478 100644 --- a/ProcessMaker/Enums/ScriptExecutorType.php +++ b/ProcessMaker/Enums/ScriptExecutorType.php @@ -7,4 +7,10 @@ enum ScriptExecutorType:string case System = 'system'; case Custom = 'custom'; case Duplicate = 'duplicate'; + case Realtime = 'realtime'; + + public function isCustomOrRealtime(): bool + { + return $this === self::Custom || $this === self::Realtime; + } } diff --git a/ProcessMaker/Events/ScriptResponseEvent.php b/ProcessMaker/Events/ScriptResponseEvent.php index f44e8fc40a..2837e11594 100644 --- a/ProcessMaker/Events/ScriptResponseEvent.php +++ b/ProcessMaker/Events/ScriptResponseEvent.php @@ -25,18 +25,21 @@ class ScriptResponseEvent implements ShouldBroadcastNow public $nonce; + public $duration; + /** * Create a new event instance. * * @return void */ - public function __construct(User $user, $status, array $response, $watcher = null, $nonce = null) + public function __construct(User $user, $status, array $response, $watcher = null, $nonce = null, $duration = null) { $this->userId = $user->id; $this->status = $status; $this->response = $response; $this->watcher = $watcher; $this->nonce = $nonce; + $this->duration = $duration; } /** @@ -75,6 +78,7 @@ public function broadcastWith() 'watcher' => $this->watcher, 'response' => $response, 'nonce' => $this->nonce, + 'duration' => $this->duration, ]; } } diff --git a/ProcessMaker/Http/Controllers/Admin/ScriptExecutorController.php b/ProcessMaker/Http/Controllers/Admin/ScriptExecutorController.php index 230458ad89..5dfe1df1fd 100644 --- a/ProcessMaker/Http/Controllers/Admin/ScriptExecutorController.php +++ b/ProcessMaker/Http/Controllers/Admin/ScriptExecutorController.php @@ -8,15 +8,20 @@ class ScriptExecutorController extends Controller { - public function index(Request $request) + public function index(Request $request, ScriptMicroserviceService $service) { if (!config('app.custom_executors')) { abort(404); } + $scriptMicroserviceEnabled = config('script-runner-microservice.enabled'); + return view('admin.script-executors.index', [ - 'script_microservice_enabled' => config('script-runner-microservice.enabled'), + 'script_microservice_enabled' => $scriptMicroserviceEnabled, + 'script_microservice_tenant_id' => $scriptMicroserviceEnabled + ? $service->getInstanceUuid() + : null, ]); } } diff --git a/ProcessMaker/Http/Controllers/Api/ScriptController.php b/ProcessMaker/Http/Controllers/Api/ScriptController.php index e3cca93814..80d1e3b78e 100644 --- a/ProcessMaker/Http/Controllers/Api/ScriptController.php +++ b/ProcessMaker/Http/Controllers/Api/ScriptController.php @@ -4,6 +4,7 @@ use Illuminate\Http\Request; use Illuminate\Support\Facades\Cache; +use ProcessMaker\Enums\ScriptExecutorType; use ProcessMaker\Events\ScriptCreated; use ProcessMaker\Events\ScriptDeleted; use ProcessMaker\Events\ScriptDuplicated; @@ -178,7 +179,7 @@ public function index(Request $request) * * @OA\Response( * response=200, - * description="success if the script was queued", + * description="The queued status or realtime execution response", * ), * ), * ) @@ -190,7 +191,12 @@ public function preview(Request $request, Script $script) $code = $request->get('code'); $nonce = $request->get('nonce'); - TestScript::dispatch($script, $request->user(), $code, $data, $config, $nonce)->onQueue('bpmn'); + $job = TestScript::dispatch($script, $request->user(), $code, $data, $config, $nonce); + if ($script->scriptExecutor->type === ScriptExecutorType::Realtime) { + $job->onConnection('redis-realtime')->onQueue('realtime'); + } else { + $job->onQueue('bpmn'); + } return ['status' => 'success']; } @@ -249,7 +255,12 @@ public function execute(Request $request, ...$scriptKey) if ($request->get('sync') === true) { return (new ExecuteScript($script, $request->user(), $code, $data, $watcher, $config, true))->handle(); } else { - ExecuteScript::dispatch($script, $request->user(), $code, $data, $watcher, $config)->onQueue('bpmn'); + $job = ExecuteScript::dispatch($script, $request->user(), $code, $data, $watcher, $config); + if ($script->scriptExecutor->type === ScriptExecutorType::Realtime) { + $job->onConnection('redis-realtime')->onQueue('realtime'); + } else { + $job->onQueue('bpmn'); + } } return ['status' => 'success', 'key' => $watcher]; diff --git a/ProcessMaker/Http/Controllers/Api/ScriptExecutorController.php b/ProcessMaker/Http/Controllers/Api/ScriptExecutorController.php index 7627c9d7a4..80bc624891 100644 --- a/ProcessMaker/Http/Controllers/Api/ScriptExecutorController.php +++ b/ProcessMaker/Http/Controllers/Api/ScriptExecutorController.php @@ -4,6 +4,7 @@ use Illuminate\Auth\Access\AuthorizationException; use Illuminate\Database\Eloquent\ModelNotFoundException; +use Illuminate\Http\Client\RequestException; use Illuminate\Http\Request; use Illuminate\Validation\ValidationException; use ProcessMaker\Enums\ScriptExecutorType; @@ -124,7 +125,15 @@ public function store(Request $request, ScriptMicroserviceService $service) ScriptExecutorCreated::dispatch($scriptExecutor->getAttributes()); BuildScriptExecutor::dispatch($scriptExecutor->id, $request->user()->id); } else { - $service->createCustomExecutor($scriptExecutor); + try { + $service->createCustomExecutor($scriptExecutor); + } catch (RequestException $e) { + // The remote executor was rejected, so keeping the local record would leave + // an executor that can never be built. + $scriptExecutor->delete(); + + $this->throwMicroserviceError($e); + } } return ['status' => 'started', 'uuid' => $scriptExecutor->uuid, 'id' => $scriptExecutor->id]; @@ -188,8 +197,12 @@ public function update(Request $request, ScriptExecutor $scriptExecutor, ScriptM $request->only($scriptExecutor->getFillable()) ); - if (config('script-runner-microservice.enabled') && $scriptExecutor->type == ScriptExecutorType::Custom) { - $service->updateCustomExecutor($scriptExecutor); + if (config('script-runner-microservice.enabled') && $scriptExecutor->type?->isCustomOrRealtime()) { + try { + $service->updateCustomExecutor($scriptExecutor); + } catch (RequestException $e) { + $this->throwMicroserviceError($e); + } } else { if (!empty($scriptExecutor->getChanges())) { ScriptExecutorUpdated::dispatch($scriptExecutor->id, $original, $scriptExecutor->getChanges()); @@ -274,6 +287,28 @@ public function delete(Request $request, ScriptExecutor $scriptExecutor, ScriptM return ['status' => 'done']; } + /** + * Report a script microservice rejection on the form, since the build itself + * is only reported over the websocket channel. + * + * @param RequestException $e Failed script microservice response + * + * @return never + * + * @throws ValidationException|RequestException + */ + private function throwMicroserviceError(RequestException $e): never + { + if (!$e->response->clientError()) { + throw $e; + } + + $detail = $e->response->json('detail'); + $message = is_string($detail) && $detail !== '' ? $detail : $e->getMessage(); + + throw ValidationException::withMessages(['language' => [$message]]); + } + private function checkAuth($request) { if (!config('app.custom_executors')) { @@ -381,11 +416,30 @@ public function availableLanguages() $languages[] = [ 'value' => $key, 'text' => $config['name'], + 'language' => $key, + 'realtime' => false, 'initDockerfile' => ScriptExecutor::initDockerfile($key), + 'configExample' => '', ]; } } + foreach (['php', 'python', 'javascript'] as $language) { + $dockerfilePath = resource_path("script-executors/realtime/{$language}.Dockerfile"); + if (!file_exists($dockerfilePath)) { + continue; + } + $label = $language === 'javascript' ? 'nodejs (realtime)' : "{$language} (realtime)"; + $languages[] = [ + 'value' => "{$language}-realtime", + 'text' => $label, + 'language' => $language, + 'realtime' => true, + 'initDockerfile' => '', + 'configExample' => file_get_contents($dockerfilePath), + ]; + } + return ['languages' => $languages]; } } diff --git a/ProcessMaker/Jobs/ErrorHandling.php b/ProcessMaker/Jobs/ErrorHandling.php index 63c3161040..b6cddf2680 100644 --- a/ProcessMaker/Jobs/ErrorHandling.php +++ b/ProcessMaker/Jobs/ErrorHandling.php @@ -83,7 +83,8 @@ private function requeue($job) ); } $newJob->delay($this->retryWaitTime()); - $newJob->onQueue('bpmn'); + $newJob->onConnection($job->connection); + $newJob->onQueue($job->queue ?? 'bpmn'); dispatch($newJob); } diff --git a/ProcessMaker/Jobs/TestScript.php b/ProcessMaker/Jobs/TestScript.php index 4a5050d2d0..a152c9bab0 100644 --- a/ProcessMaker/Jobs/TestScript.php +++ b/ProcessMaker/Jobs/TestScript.php @@ -7,10 +7,12 @@ use Illuminate\Foundation\Bus\Dispatchable; use Illuminate\Queue\InteractsWithQueue; use Illuminate\Queue\SerializesModels; +use Illuminate\Support\Facades\Log; use ProcessMaker\Enums\ScriptExecutorType; use ProcessMaker\Events\ScriptResponseEvent; use ProcessMaker\Models\Script; use ProcessMaker\Models\User; +use ProcessMaker\Services\ScriptMicroserviceService; use Throwable; class TestScript implements ShouldQueue @@ -35,8 +37,8 @@ class TestScript implements ShouldQueue /** * Create a new job instance to execute a script. * - * @param ProcessMaker\Models\Script $script - * @param ProcessMaker\Models\User $current_user + * @param Script $script + * @param User $current_user * @param string $code * @param array $data * @param array $configuration @@ -58,24 +60,31 @@ public function __construct(Script $script, User $current_user, $code, array $da */ public function handle() { + $startTime = microtime(true); try { // Just set the code but do not save the object (preview only) $this->script->code = $this->code; $metadata = [ 'nonce' => $this->nonce, 'current_user' => $this->current_user?->id, + 'start_time' => microtime(true), ]; $response = $this->script->runScript($this->data, $this->configuration, '', null, 0, $metadata); - \Log::debug('Response from runScript: ' . print_r($response, true)); + Log::debug('Response from runScript: ' . print_r($response, true)); - if (!config('script-runner-microservice.enabled')) { - $this->sendResponse(200, $response); + if ($this->script->scriptExecutor->type === ScriptExecutorType::Realtime) { + // Realtime executors return the microservice payload synchronously rather than + // posting back to the callback endpoint, so format it the same way here. + $formatted = (new ScriptMicroserviceService())->formatPreviewResponse($response); + $this->sendResponse($formatted['status'], $formatted['output'], (microtime(true) - $startTime)); + } elseif (!config('script-runner-microservice.enabled')) { + $this->sendResponse(200, $response, (microtime(true) - $startTime)); } } catch (Throwable $exception) { $this->sendResponse(500, [ 'exception' => get_class($exception), 'message' => $exception->getMessage(), - ]); + ], (microtime(true) - $startTime)); } } @@ -85,8 +94,8 @@ public function handle() * @param int $status * @param array $response */ - private function sendResponse($status, array $response) + private function sendResponse($status, array $response, float $duration) { - event(new ScriptResponseEvent($this->current_user, $status, $response, null, $this->nonce)); + event(new ScriptResponseEvent($this->current_user, $status, $response, null, $this->nonce, $duration)); } } diff --git a/ProcessMaker/Models/Script.php b/ProcessMaker/Models/Script.php index f01c7f4492..3853001bfd 100644 --- a/ProcessMaker/Models/Script.php +++ b/ProcessMaker/Models/Script.php @@ -179,6 +179,7 @@ public function runScript(array $data, array $config, $tokenId = '', $timeout = if (!$user) { throw new ConfigurationException('A user is required to run scripts'); } + $metadata['start_time'] = microtime(true); return $runner->run($this->code, $data, $config, $timeout, $user, $sync, $metadata); } diff --git a/ProcessMaker/Models/ScriptExecutor.php b/ProcessMaker/Models/ScriptExecutor.php index 6e5438f279..408b9e457d 100644 --- a/ProcessMaker/Models/ScriptExecutor.php +++ b/ProcessMaker/Models/ScriptExecutor.php @@ -32,6 +32,7 @@ * @OA\Property(property="language", type="string"), * @OA\Property(property="config", type="string"), * @OA\Property(property="is_system", type="boolean"), + * @OA\Property(property="type", type="string"), * ), * @OA\Schema( * schema="scriptExecutors", @@ -177,6 +178,12 @@ public static function rules($existing = null) }); } + // Realtime executors may use php/python/javascript even if a package is not installed + $allowedLanguages = array_values(array_unique(array_merge( + $allowedLanguages, + ['php', 'python', 'javascript'] + ))); + return [ 'title' => 'required', 'language' => [ diff --git a/ProcessMaker/Nayra/Managers/WorkflowManagerDefault.php b/ProcessMaker/Nayra/Managers/WorkflowManagerDefault.php index 11a8d3007f..cb521091a7 100644 --- a/ProcessMaker/Nayra/Managers/WorkflowManagerDefault.php +++ b/ProcessMaker/Nayra/Managers/WorkflowManagerDefault.php @@ -8,6 +8,7 @@ use ProcessMaker\BpmnEngine; use ProcessMaker\Contracts\ServiceTaskImplementationInterface; use ProcessMaker\Contracts\WorkflowManagerInterface; +use ProcessMaker\Enums\ScriptExecutorType; use ProcessMaker\Jobs\BoundaryEvent; use ProcessMaker\Jobs\CallProcess; use ProcessMaker\Jobs\CatchEvent; @@ -23,6 +24,7 @@ use ProcessMaker\Models\Process as Definitions; use ProcessMaker\Models\ProcessRequest; use ProcessMaker\Models\ProcessRequestToken as Token; +use ProcessMaker\Models\Script; use ProcessMaker\Nayra\Contracts\Bpmn\BoundaryEventInterface; use ProcessMaker\Nayra\Contracts\Bpmn\EntityInterface; use ProcessMaker\Nayra\Contracts\Bpmn\EventDefinitionInterface; @@ -231,7 +233,14 @@ public function runScripTask(ScriptTaskInterface $scriptTask, Token $token) $instance = $token->processRequest; $process = $instance->process; - RunScriptTask::dispatch($process, $instance, $token, [])->onQueue('bpmn'); + $job = new RunScriptTask($process, $instance, $token, []); + $script = Script::find($scriptTask->getProperty('scriptRef')); + + if ($script?->scriptExecutor?->type === ScriptExecutorType::Realtime) { + $job->onConnection('redis-realtime')->onQueue('realtime'); + } + + dispatch($job); } /** diff --git a/ProcessMaker/ScriptRunners/ScriptMicroserviceRunner.php b/ProcessMaker/ScriptRunners/ScriptMicroserviceRunner.php index dc4d5f3a96..a4c603d117 100644 --- a/ProcessMaker/ScriptRunners/ScriptMicroserviceRunner.php +++ b/ProcessMaker/ScriptRunners/ScriptMicroserviceRunner.php @@ -40,7 +40,7 @@ public function run($code, array $data, array $config, $timeout, $user, $sync, $ $scriptRunner = $this->service->getScriptRunner( $this->language, $this->script->scriptExecutor->uuid, - $this->script->scriptExecutor->type === ScriptExecutorType::Custom + $this->script->scriptExecutor->type?->isCustomOrRealtime() ); if (!$scriptRunner) { @@ -62,7 +62,8 @@ public function run($code, array $data, array $config, $timeout, $user, $sync, $ 'callback_token' => $environmentVariables['API_TOKEN'], 'debug' => true, 'timeout' => $timeout, - 'sync' => $sync, + // Realtime executors always run synchronously (service ignores sync too) + 'sync' => $this->script->scriptExecutor->type === ScriptExecutorType::Realtime ? true : $sync, ]; Log::debug('Payload: ' . print_r($payload, true)); diff --git a/ProcessMaker/Services/ScriptMicroserviceService.php b/ProcessMaker/Services/ScriptMicroserviceService.php index b3a8a114ce..61098bcaee 100644 --- a/ProcessMaker/Services/ScriptMicroserviceService.php +++ b/ProcessMaker/Services/ScriptMicroserviceService.php @@ -9,6 +9,7 @@ use Illuminate\Support\Facades\Cache; use Illuminate\Support\Facades\Http; use Illuminate\Support\Facades\Log; +use ProcessMaker\Enums\ScriptExecutorType; use ProcessMaker\Events\ScriptResponseEvent; use ProcessMaker\Exception\ScriptException; use ProcessMaker\Jobs\CompleteActivity; @@ -87,6 +88,7 @@ public function createCustomExecutor(ScriptExecutor $scriptExecutor) 'language' => strtolower($scriptExecutor->language), 'version' => config('script-runner-microservice.version'), 'config' => $scriptExecutor->config, + 'realtime' => $scriptExecutor->type === ScriptExecutorType::Realtime, ]; Log::debug('Payload: ', $payload); @@ -114,6 +116,7 @@ public function updateCustomExecutor(ScriptExecutor $scriptExecutor) 'language' => strtolower($scriptExecutor->language), 'version' => config('script-runner-microservice.version'), 'config' => $scriptExecutor->config, + 'realtime' => $scriptExecutor->type === ScriptExecutorType::Realtime, ]; Log::debug('Payload: ', $payload); @@ -195,6 +198,7 @@ public function getScriptRunner(string $language, string $executorUuid, bool $cu if (Cache::has($cacheKey)) { Log::debug('Cache hit for script runner', ['cacheKey' => $cacheKey]); + return Cache::get($cacheKey); } @@ -211,7 +215,9 @@ public function getScriptRunner(string $language, string $executorUuid, bool $cu return isset($item['language'], $item['id']) && $item['language'] === $language && $item['id'] === $executorUuid; })->first(); - if (!empty($result)) Cache::put($cacheKey, $result, now()->addHour()); + if (!empty($result)) { + Cache::put($cacheKey, $result, now()->addHour()); + } return $result; } @@ -235,6 +241,9 @@ public function handle(Request $request): void { $response = $request->all(); Log::debug('Response microservice executor: ' . print_r($response, true)); + + $duration = isset($response['metadata']['start_time']) ? (microtime(true) - $response['metadata']['start_time']) : null; + // If the call is from preview if (!empty($response['metadata']['nonce'])) { $formattedResponse = $this->formatPreviewResponse($response); @@ -243,7 +252,8 @@ public function handle(Request $request): void $formattedResponse['status'], $formattedResponse['output'], null, - $response['metadata']['nonce'])); + $response['metadata']['nonce'], + $duration)); } if (!empty($response['metadata']['script_task'])) { $script = Script::find($response['metadata']['script_task']['script_id']); @@ -262,10 +272,10 @@ public function handle(Request $request): void * @param array $response * @return array{status: int, output: array} */ - private function formatPreviewResponse(array $response): array + public function formatPreviewResponse(array $response): array { // Simple status determination: success = 200, others = 500 - $status = $response['status'] === 'success' ? 200 : 500; + $status = ($response['status'] ?? '') === 'success' ? 200 : 500; return [ 'status' => $status, @@ -284,7 +294,7 @@ private function formatPreviewOutput(array $response): array $output = $response; if (($response['status'] ?? '') === 'success') { $output = ['output' => $response['output']]; - } elseif ($response['status'] === 'error') { + } elseif (($response['status'] ?? '') === 'error') { $output = [ 'exception' => $response['exception'] ?? ScriptException::class, 'message' => $response['error'], diff --git a/config/horizon.php b/config/horizon.php index af57f786e1..4569dad6de 100644 --- a/config/horizon.php +++ b/config/horizon.php @@ -163,6 +163,16 @@ 'maxProcesses' => env('PM4_HORIZON_SUPERVISOR_1_MAX_PROCESSES', 1), 'memory' => env('PM4_HORIZON_WORKER_MEMORY_LIMIT', 512), ], + 'supervisor-realtime' => [ + 'connection' => 'redis-realtime', + 'queue' => ['realtime'], + 'balance' => env('PM4_HORIZON_SUPERVISOR_REALTIME_BALANCE', 'auto'), // auto: dynamically adjust the number of processes based on the load + 'tries' => env('PM4_HORIZON_SUPERVISOR_REALTIME_TRIES', 3), + 'timeout' => env('PM4_HORIZON_SUPERVISOR_REALTIME_TIMEOUT', 600), + 'minProcesses' => env('PM4_HORIZON_SUPERVISOR_REALTIME_MIN_PROCESSES', 2), + 'maxProcesses' => env('PM4_HORIZON_SUPERVISOR_REALTIME_MAX_PROCESSES', 20), + 'memory' => env('PM4_HORIZON_WORKER_MEMORY_LIMIT', 512), + ], ], 'local' => [ @@ -187,6 +197,16 @@ 'maxProcesses' => env('PM4_HORIZON_SUPERVISOR_1_MAX_PROCESSES', 1), 'memory' => env('PM4_HORIZON_WORKER_MEMORY_LIMIT', 512), ], + 'supervisor-realtime' => [ + 'connection' => 'redis-realtime', + 'queue' => ['realtime'], + 'balance' => env('PM4_HORIZON_SUPERVISOR_REALTIME_BALANCE', 'auto'), // auto: dynamically adjust the number of processes based on the load + 'tries' => env('PM4_HORIZON_SUPERVISOR_REALTIME_TRIES', 3), + 'timeout' => env('PM4_HORIZON_SUPERVISOR_REALTIME_TIMEOUT', 600), + 'minProcesses' => env('PM4_HORIZON_SUPERVISOR_REALTIME_MIN_PROCESSES', 2), + 'maxProcesses' => env('PM4_HORIZON_SUPERVISOR_REALTIME_MAX_PROCESSES', 20), + 'memory' => env('PM4_HORIZON_WORKER_MEMORY_LIMIT', 512), + ], ], 'staging' => [ @@ -211,6 +231,16 @@ 'maxProcesses' => env('PM4_HORIZON_SUPERVISOR_1_MAX_PROCESSES', 1), 'memory' => env('PM4_HORIZON_WORKER_MEMORY_LIMIT', 512), ], + 'supervisor-realtime' => [ + 'connection' => 'redis-realtime', + 'queue' => ['realtime'], + 'balance' => env('PM4_HORIZON_SUPERVISOR_REALTIME_BALANCE', 'auto'), + 'tries' => env('PM4_HORIZON_SUPERVISOR_REALTIME_TRIES', 3), + 'timeout' => env('PM4_HORIZON_SUPERVISOR_REALTIME_TIMEOUT', 600), + 'minProcesses' => env('PM4_HORIZON_SUPERVISOR_REALTIME_MIN_PROCESSES', 2), + 'maxProcesses' => env('PM4_HORIZON_SUPERVISOR_REALTIME_MAX_PROCESSES', 20), + 'memory' => env('PM4_HORIZON_WORKER_MEMORY_LIMIT', 512), + ], ], ], ]; diff --git a/config/queue.php b/config/queue.php index c5029a5baa..6bda789295 100644 --- a/config/queue.php +++ b/config/queue.php @@ -71,6 +71,15 @@ 'after_commit' => false, ], + 'redis-realtime' => [ + 'driver' => 'redis', + 'connection' => 'default', + 'queue' => 'realtime', + 'retry_after' => 86400, + 'block_for' => 10, + 'after_commit' => false, + ], + ], /* diff --git a/resources/js/admin/script-executors/ScriptExecutors.vue b/resources/js/admin/script-executors/ScriptExecutors.vue index 63970dc86c..1de44521cb 100644 --- a/resources/js/admin/script-executors/ScriptExecutors.vue +++ b/resources/js/admin/script-executors/ScriptExecutors.vue @@ -32,7 +32,7 @@ {{ props.rowData.title }}