PHP code example of tbessenreither / quickfork

1. Go to this page and download the library: Download tbessenreither/quickfork library. Choose the download type require.

2. Extract the ZIP file and open the index.php.

3. Add this code to the index.php.
    
        
<?php
require_once('vendor/autoload.php');

/* Start to develop here. Best regards https://php-download.com/ */

    

tbessenreither / quickfork example snippets



use Tbessenreither\Quickfork\Objects\Task;
use Tbessenreither\Quickfork\Quickfork;

class FastparallelCommand
{

    protected function execute(): int
    {
        $io = new SymfonyStyle($input, $output);

        $io->title('Running QuickFork Command');

        $quickfork = new Quickfork();

        try {
            $tasks = [];
            for ($i = 0; $i < 4; $i++) {
                $task = new Task(
                    callable: $this->workerThread(...),
                    arguments: [
                        'number' => $i + 1,
                        'index' => $i,
                    ],
                );
                $tasks[] = $task;
            }

            $results = $quickfork->runTasksInThreads($tasks, maxConcurrent: 5);

            foreach ($tasks as $task) {
                $taskId = $task->getId();
                $taskResult = $results[$taskId] ?? null;
                echo "Task ID: {$taskId}, Result:\n";
                if ($taskResult->hasError()) {
                    echo "Error: " . $taskResult->getError()->getMessage() . "\n";
                } else {
                    echo "Result: " . $taskResult->getResult() . "\n";
                }
                echo "-----------------------------\n";
            }

            $io->success('QuickFork Command executed successfully.');
        } catch (Throwable $e) {
            $io->error('An error occurred while executing QuickFork Command: ' . $e->getMessage());

            return Command::FAILURE;
        }


        return Command::SUCCESS;
    }

    private function workerThread(Task $task, ?int $index = null, ?int $number = null): string
    {
        //sleep(rand(1, 3));
        if ($number === 3) {
            throw new RuntimeException('Simulated critical error in worker thread.');
        }
        return "Worker thread is executing task {$number}.";
    }

}

class QuickforkCommand extends Command
{
    private int $numberOfTasks = 5;

    protected function execute(InputInterface $input, OutputInterface $output): int
    {
        $io = new SymfonyStyle($input, $output);

        $io->title('Running QuickFork Command');

        $quickfork = new Quickfork();

        $task = new Task(
            callable: $this->taskHandler(...),
            arguments: [],
        );
        $quickfork->runTask($task);

        while ($task->isRunning()) {
            $messages = $task->getSocket()->getMessages();
            $this->handleParentMessages($messages, $task->getSocket());
            usleep(100 * 1000); // Sleep for 100ms to prevent busy waiting
        }

        $finalMessages = $task->getSocket()->getMessages();
        $this->handleParentMessages($finalMessages, $task->getSocket());

        return Command::SUCCESS;
    }

    /**
     * @param Message[] $messages
     */
    private function handleParentMessages(array $messages, Socket $socket): void
    {
        foreach ($messages as $message) {
            echo "Received message: Topic: {$message->getTopic()}\n";

            if ($message->getTopic() === 'fork_start') {
                echo "Task started with ID: {$message->getForkId()}\n";
            } elseif ($message->getTopic() === 'fork_complete') {
                echo "Task with ID {$message->getForkId()} completed.\n";
            } elseif ($message->getTopic() === 'fork_output') {
                echo "Output: {$message->getContent()}\n";
            } elseif ($message->getTopic() === 'fork_result') {
                echo "Result:\n";
                var_dump($message->getContent());
                echo "-----------------------------\n";
            } elseif ($message->getTopic() === 'fork_error') {
                echo "Error:\n";
                var_dump($message->getContent());
                echo "-----------------------------\n";
            } elseif ($message->getTopic() === 'ready_for_task') {
                echo "Worker thread is ready for task. Sending out new one.\n";
                if ($this->numberOfTasks > 0) {
                    $socket->send(new Message(
                        topic: 'new_task',
                        content: [
                            'sum' => [rand(1, 100), rand(1, 100)],
                        ]
                    ));
                    $this->numberOfTasks--;
                    if ($this->numberOfTasks <= 0) {
                        $socket->send(new Message(
                            topic: 'shutdown',
                        ));
                    }
                }
            } elseif ($message->getTopic() === 'thread_result') {
                echo "Thread Result:\n";
                var_dump($message->getContent());
                echo "-----------------------------\n";
            }
        }
    }

    private function taskHandler(Task $task): string
    {
        $task->getSocket()->send(new Message(
            topic: 'ready_for_task',
        ));

        $active = true;
        while ($active) {
            $messages = $task->getSocket()->getMessages();
            foreach ($messages as $message) {
                if ($message->getTopic() === 'new_task') {
                    $content = $message->getContent();
                    $sum = array_sum($content['sum']);
                    echo "doing more stupid calculations for the boss...\n";
                    $task->getSocket()->send(new Message(
                        topic: 'thread_result',
                        content: "The sum of {$content['sum'][0]} and {$content['sum'][1]} is {$sum}.",
                        replyTo: $message->getId(),
                    ));
                    $task->getSocket()->send(new Message(
                        topic: 'ready_for_task',
                    ));
                } elseif ($message->getTopic() === 'shutdown') {
                    $active = false;
                    return "Task is shutting down.";
                }
            }
            usleep(100 * 1000); // Sleep for 100ms to prevent busy waiting
        }
        // Simulate some work
        sleep(rand(1, 3));
        return "Task with ID {$task->getId()} completed.";
    }
}