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.";
}
}