1. Go to this page and download the library: Download yangusik/thrun 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/ */
yangusik / thrun example snippets
use Thrun\Envelope\Envelope;
use Thrun\Supervisor\Supervisor;
use Thrun\Supervisor\SupervisorOptions;
use Thrun\Transport\InMemory\InMemoryTransport;
use Thrun\Worker\Worker;
use Thrun\Worker\WorkerOptions;
$transport = new InMemoryTransport();
// Send messages
$transport->send(Envelope::wrap(new SendEmailMessage('[email protected]', 'Hello')));
$supervisor = new Supervisor(
workerFactory: fn() => new Worker(
transport: $transport,
handlers: [
SendEmailMessage::class => fn(SendEmailMessage $m) => mail($m->to, $m->subject, '...'),
],
options: new WorkerOptions(threads: 2, concurrency: 4),
),
options: new SupervisorOptions(),
);
$supervisor->run();
use Thrun\Envelope\Stamp\TimeoutStamp;
$transport->send(Envelope::wrap(
new SlowApiCall(),
new TimeoutStamp(5000), // 5 seconds
));
use Thrun\Envelope\Stamp\RetryStamp;
use Thrun\Envelope\Stamp\MessageIdStamp;
$transport->send(Envelope::wrap(
new SendEmailMessage('[email protected]', 'Hello'),
new RetryStamp(backoff: [1000, 2000, 4000], maxAttempts: 3),
new MessageIdStamp('msg-42'),
));
use Thrun\Worker\Metrics\InMemoryMetrics;
$metrics = new InMemoryMetrics();
$worker = new Worker(transport: $transport, handlers: $handlers, metrics: $metrics);
// $metrics->processed, $metrics->failed, $metrics->retried, $metrics->timedOut
// $metrics->averageTime()
use Thrun\Envelope\Stamp\PartitionStamp;
use Thrun\Transport\Policy\MaxConcurrencyPolicy;
use Thrun\Transport\PolicyAwareReceiver;
$receiver = new PolicyAwareReceiver(
inner: $transport,
policy: new MaxConcurrencyPolicy(maxPerPartition: 5),
);
$transport->send(Envelope::wrap($msg, new PartitionStamp('tenant-42')));
use Thrun\Worker\Acknowledger;
SendEmailMessage::class => function (SendEmailMessage $m, Acknowledger $ack) {
if ($m->to === '[email protected]') {
$ack->fail(new \RuntimeException('Blocked'));
return;
}
// process...
$ack->ack();
}
use Thrun\Middleware\CatchMessageMiddleware;
$worker = new Worker(
transport: $transport,
handlers: $handlers,
options: new WorkerOptions(middleware: [new CatchMessageMiddleware()]),
);
use Thrun\Transport\Redis\RedisTransport;
use Thrun\Transport\Redis\RedisConnection;
use Thrun\Serialization\JsonSerializer;
$redis = new \Redis([
'pool' => [
'enabled' => true,
'min' => 0,
'max' => 1,
'mux' => 0,
],
]);
$redis->connect('redis', 6379);
$transport = new RedisTransport(
connection: new RedisConnection($redis, 'thrun:queue'),
serializer: new JsonSerializer(),
queue: 'emails',
);
$failureTransport = new InMemoryTransport(); // or RedisTransport
$worker = new Worker(
transport: $transport,
handlers: $handlers,
failureTransport: $failureTransport,
);