Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion composer.json
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
{
"name": "blitz-php/queue",
"description": "Gestionnaire de file d'attente pour BlitzPHP",
"keywords": ["blitz-php", "blitz php", "queue", "worker", "database", "redis", "predis", "file d'attente" ],
"keywords": ["blitz-php", "blitz php", "queue", "worker", "database", "redis", "predis" ],
"homepage": "https://github.com/blitz-php/queue",
"license": "MIT",
"type": "library",
Expand Down
25 changes: 11 additions & 14 deletions src/Commands/Work.php
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,6 @@

namespace BlitzPHP\Queue\Commands;

use BlitzPHP\Cache\Handlers\BaseHandler;
use BlitzPHP\Cli\Console\Command;
use BlitzPHP\Cli\Console\Console;
use BlitzPHP\Contracts\Cache\CacheInterface;
Expand Down Expand Up @@ -119,8 +118,6 @@ public function __construct(protected ContainerInterface $container, protected C
$this->worker = service('worker');
$this->cache = $container->get(CacheInterface::class);
$this->events = $container->get(EventManagerInterface::class);

BaseHandler::setReservedCharacters(str_replace(':', '', config('cache.reserved_characters')));
}

/**
Expand Down Expand Up @@ -162,7 +159,7 @@ public function execute(array $params)
protected function runWorker(string $connection, string $queue): ?int
{
return $this->worker
->setName($this->option('name'))
->setName($this->option('name', 'default'))
->setCache($this->cache)
->{$this->option('once') ? 'runNextJob' : 'daemon'}(
$connection,
Expand All @@ -177,17 +174,17 @@ protected function runWorker(string $connection, string $queue): ?int
protected function gatherWorkerOptions(): WorkerOptions
{
return new WorkerOptions(
$this->option('name'),
max($this->option('backoff'), $this->option('delay')),
$this->option('memory'),
$this->option('timeout'),
$this->option('sleep'),
$this->option('tries'),
$this->option('name', 'default'),
max($this->option('backoff', 0), $this->option('delay', 0)),
$this->option('memory', 128),
$this->option('timeout', 60),
$this->option('sleep', 3),
$this->option('tries', 1),
$this->option('force', false),
$this->option('stop-when-empty', false),
$this->option('max-jobs'),
$this->option('max-time'),
$this->option('rest'),
$this->option('max-jobs', 0),
$this->option('max-time', 0),
$this->option('rest', 0),
);
}

Expand Down Expand Up @@ -351,7 +348,7 @@ protected function getQueue(string $connection): string
/**
* Indique si l'application est en maintenance (et si le worker doit s'arrêter).
*/
protected function downForMaintenance(): false
protected function downForMaintenance(): bool
{
return $this->option('force')
? false
Expand Down
4 changes: 3 additions & 1 deletion src/Config/Services.php
Original file line number Diff line number Diff line change
Expand Up @@ -76,7 +76,9 @@ public static function worker(bool $shared = true): Worker
}
}

memory_reset_peak_usage();
if (function_exists('memory_reset_peak_usage')) {
memory_reset_peak_usage();
}
};

return static::$instances[Worker::class] = new Worker(
Expand Down
53 changes: 31 additions & 22 deletions src/Config/queue.php
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,9 @@
*/

use BlitzPHP\Queue\Drivers\DatabaseDriver;
use BlitzPHP\Queue\Drivers\FailoverDriver;
use BlitzPHP\Queue\Drivers\NullDriver;
use BlitzPHP\Queue\Drivers\SyncDriver;

/**
* Configuration du composant de files d'attente (queue).
Expand Down Expand Up @@ -49,7 +52,7 @@
* Groupe / nom de connexion base de données BlitzPHP à utiliser
* pour lire et écrire les jobs. Variable : `queue.database.group`.
*/
'group' => env('queue.database.group', 'default'),
'connection' => env('queue.database.group', 'default'),

/**
* Si `true`, réutilise une connexion partagée du gestionnaire de
Expand All @@ -70,23 +73,23 @@
*/
'table' => env('queue.database.table', 'queue_jobs'),

/**
* Nom de la file logique par défaut pour cette connexion
* (colonne `queue` en base). Utilisé si `queue:work` n'en précise pas.
*/
// 'queue' => 'default',
/**
* Nom de la file logique par défaut pour cette connexion
* (colonne `queue` en base). Utilisé si `queue:work` n'en précise pas.
*/
'queue' => env('queue.defaultQueue', 'default'),

/**
* Délai en secondes au-delà duquel un job réservé est considéré
* comme expiré et peut être repris par un autre worker.
*/
// 'retry_after' => 60,
/**
* Délai en secondes au-delà duquel un job réservé est considéré
* comme expiré et peut être repris par un autre worker.
*/
'retry_after' => (int) env('queue.retryAfter', 90),

/**
* Si `true`, n'envoie le job qu'après le commit des transactions
* de base de données en cours.
*/
// 'after_commit' => false,
/**
* Si `true`, n'envoie le job qu'après le commit des transactions
* de base de données en cours.
*/
'after_commit' => false,
],

/**
Expand Down Expand Up @@ -215,9 +218,12 @@
* Pilote SQL : table `queue_jobs` (ou celle configurée).
*/
'database' => DatabaseDriver::class,
// 'redis' => \BlitzPHP\Queue\Drivers\Redis::class,
// 'predis' => \BlitzPHP\Queue\Drivers\Predis::class,
// 'rabbitmq' => \BlitzPHP\Queue\Drivers\RabbitMQ::class,
// 'redis' => \BlitzPHP\Queue\Drivers\RedisDriver::class,
// 'predis' => \BlitzPHP\Queue\Drivers\PredisDriver::class,
// 'rabbitmq' => \BlitzPHP\Queue\Drivers\RabbitMQDriver::class,
'sync' => SyncDriver::class,
'null' => NullDriver::class,
'failover' => FailoverDriver::class,
],

/**
Expand Down Expand Up @@ -249,10 +255,13 @@
'database' => env('db.connection', 'default'),

/**
* Table SQL des jobs échoués (`uuid`, `connection`, `queue`, `payload`,
* `exception`, `failed_at`).
* Table SQL des jobs échoués (`uuid`, `connection`, `queue`, `payload`, `exception`, `failed_at`).
*/
'table' => 'queue_failed_jobs',

// Pour le driver 'file'
// 'path' => storage_path('logs/failed_jobs.json'),
// 'limit' => 100,
],

/**
Expand All @@ -269,6 +278,6 @@
/**
* Nom de la table (ou identifiant de stockage) des lots de jobs.
*/
'table' => 'queue.job_batches',
'table' => 'queue_job_batches',
],
];
43 changes: 43 additions & 0 deletions src/Job.php
Original file line number Diff line number Diff line change
Expand Up @@ -11,9 +11,12 @@

namespace BlitzPHP\Queue;

use BlitzPHP\Queue\Config\Services;
use BlitzPHP\Queue\Traits\Dispatchable;
use BlitzPHP\Queue\Traits\InteractsWithQueue;
use BlitzPHP\Queue\Traits\SerializesModels;
use DateInterval;
use DateTimeInterface;

/**
* Classe de base des jobs métier destinés à la file d'attente.
Expand Down Expand Up @@ -66,4 +69,44 @@ public function queue(): string
{
return $this->queue;
}

/**
* Ajoute le job sur la queue dans la file d'attente
*/
public function push(): mixed
{
return $this->pushOn($this->queue);
}

/**
* Ajoute le job sur une queue spécifique dans la file d'attente
*/
public function pushOn(string $queue): mixed
{
return Services::queue()->pushOn($queue, $this);
}

/**
* Ajoute le job avec délai dans la file d'attente
*/
public function pushLater(DateInterval|DateTimeInterface|int $delay): mixed
{
return $this->pushLaterOn($this->queue, $delay);
}

/**
* Ajoute le job sur une queue spécifique avec délai dans la file d'attente
*/
public function pushLaterOn(string $queue, DateInterval|DateTimeInterface|int $delay): mixed
{
return Services::queue()->laterOn($queue, $delay, $this);
}

/**
* Execute le job immédiatement (synchrone)
*/
public function execute(): void
{
Services::container()->call([$this, 'handle']);
}
}
Loading