PHP code example of darkjest / deferq

1. Go to this page and download the library: Download darkjest/deferq 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/ */

    

darkjest / deferq example snippets




use DeferQ\DeferQ;
use DeferQ\Fingerprint\DefaultFingerprintGenerator;
use DeferQ\Lock\CacheLock;
use DeferQ\Queue\RedisQueueAdapter;
use DeferQ\Result\CacheResultStore;
use DeferQ\Store\CacheTaskStore;
use Predis\Client;
use Symfony\Component\Cache\Psr16Cache;
use Symfony\Component\Cache\Adapter\RedisAdapter;

$redis = new Client('tcp://127.0.0.1:6379');
$cache = new Psr16Cache(new RedisAdapter($redis));

$deferq = new DeferQ(
    queue: new RedisQueueAdapter($redis),
    taskStore: new CacheTaskStore($cache),
    resultStore: new CacheResultStore($cache),
    lock: new CacheLock($cache),
    fingerprinter: new DefaultFingerprintGenerator(),
);



use DeferQ\Handler\TaskHandlerInterface;
use DeferQ\Handler\TaskHandlerRegistry;
use DeferQ\Task\Task;

class ReportGenerateHandler implements TaskHandlerInterface
{
    public function handle(Task $task): mixed
    {
        $year = $task->params['year'];
        $format = $task->params['format'];

        // ... heavy computation ...

        return ['url' => "/reports/report-{$year}.{$format}", 'rows' => 15000];
    }
}

$handlers = new TaskHandlerRegistry();
$handlers->register('report.generate', new ReportGenerateHandler());



use DeferQ\Task\TaskStatus;

$receipt = $deferq->dispatch(
    name: 'report.generate',
    params: ['year' => 2024, 'format' => 'xlsx'],
    resultTtl: 3600,
);

match ($receipt->status) {
    TaskStatus::Completed => handleReady($receipt->result),
    TaskStatus::Running   => pollLater($receipt->taskId),
    TaskStatus::Pending   => pollLater($receipt->taskId),
    TaskStatus::Failed    => handleError($receipt->taskId),
};

$receipt = $deferq->getStatus($taskId);

if ($receipt->status === TaskStatus::Completed) {
    $result = $receipt->result;
}



use DeferQ\Worker\Worker;
use DeferQ\Worker\WorkerConfig;

$worker = new Worker(
    queue: $queueAdapter,
    handlers: $handlers,
    taskStore: $taskStore,
    resultStore: $resultStore,
    config: new WorkerConfig(
        sleepMs: 1000,
        maxJobs: 500,
        maxMemoryMb: 128,
        taskTimeoutSeconds: 300,
    ),
    logger: $psrLogger,
);

$worker->run();

// Both calls return the same task receipt (deduplication)
$receipt1 = $deferq->dispatch('report.generate', ['year' => 2024, 'format' => 'xlsx']);
$receipt2 = $deferq->dispatch('report.generate', ['format' => 'xlsx', 'year' => 2024]);

// $receipt1->taskId === $receipt2->taskId



use DeferQ\Callback\CallbackInterface;
use DeferQ\Task\Task;

class WebSocketNotifier implements CallbackInterface
{
    public function __construct(private WebSocketServer $ws) {}

    public function __invoke(Task $task, mixed $result): void
    {
        $this->ws->send($task->id, json_encode($result));
    }
}

$receipt = $deferq->dispatch(
    name: 'report.generate',
    params: ['year' => 2024],
    callback: new WebSocketNotifier($ws),
);



use DeferQ\Callback\CallbackChain;

$receipt = $deferq->dispatch(
    name: 'report.generate',
    params: ['year' => 2024],
    callback: new CallbackChain(
        new WebSocketNotifier($ws),
        new EmailNotifier($mailer),
        new MetricsRecorder($metrics),
    ),
);



use DeferQ\Queue\QueueAdapterInterface;
use DeferQ\Task\Task;

class DatabaseQueueAdapter implements QueueAdapterInterface
{
    public function __construct(private PDO $pdo) {}

    public function push(Task $task): void
    {
        $stmt = $this->pdo->prepare(
            'INSERT INTO deferq_queue (payload, created_at) VALUES (?, NOW())'
        );
        $stmt->execute([json_encode($task->toArray())]);
    }

    public function pop(int $timeoutSeconds = 5): ?Task
    {
        // Fetch and lock the oldest unprocessed row
        $stmt = $this->pdo->prepare(
            'SELECT id, payload FROM deferq_queue
             WHERE processing = 0
             ORDER BY created_at ASC
             LIMIT 1
             FOR UPDATE SKIP LOCKED'
        );
        $stmt->execute();
        $row = $stmt->fetch(PDO::FETCH_ASSOC);

        if (!$row) {
            return null;
        }

        $this->pdo->prepare('UPDATE deferq_queue SET processing = 1 WHERE id = ?')
            ->execute([$row['id']]);

        return Task::fromArray(json_decode($row['payload'], true));
    }

    public function ack(Task $task): void
    {
        // Delete processed row
    }

    public function nack(Task $task): void
    {
        // Reset processing flag
    }
}


// worker-bootstrap.php

rRegistry;
use DeferQ\Lock\CacheLock;
use DeferQ\Queue\RedisQueueAdapter;
use DeferQ\Result\CacheResultStore;
use DeferQ\Store\CacheTaskStore;
use DeferQ\Worker\Worker;
use DeferQ\Worker\WorkerConfig;
use Predis\Client;
use Psr\Log\NullLogger;

// Your PSR-16 cache implementation
$redis = new Client('tcp://127.0.0.1:6379');
$cache = /* your PSR-16 cache backed by Redis */;

// Register handlers
$handlers = new TaskHandlerRegistry();
$handlers->register('report.generate', new ReportGenerateHandler());
$handlers->register('export.csv', new CsvExportHandler());

// CLI overrides from command-line arguments
$overrides = $GLOBALS['deferq_cli_overrides'] ?? [];

return new Worker(
    queue: new RedisQueueAdapter($redis),
    handlers: $handlers,
    taskStore: new CacheTaskStore($cache),
    resultStore: new CacheResultStore($cache),
    config: new WorkerConfig(
        sleepMs: $overrides['sleepMs'] ?? 1000,
        maxJobs: $overrides['maxJobs'] ?? 0,
        maxMemoryMb: $overrides['maxMemoryMb'] ?? 128,
        taskTimeoutSeconds: $overrides['taskTimeoutSeconds'] ?? 300,
    ),
    logger: new NullLogger(),
);
bash
composer 
bash
php bin/deferq-worker --bootstrap=worker-bootstrap.php --max-jobs=1000 --sleep=500
bash
php bin/deferq-worker --bootstrap=worker-bootstrap.php --max-jobs=1000 --sleep=500