.github
2 years ago
CannotObtainLockException.php
4 years ago
FsConnectionFactory.php
4 years ago
FsConsumer.php
4 years ago
FsContext.php
4 years ago
FsDestination.php
4 years ago
FsMessage.php
4 years ago
FsProducer.php
4 years ago
LICENSE
4 years ago
LegacyFilesystemLock.php
4 years ago
Lock.php
4 years ago
README.md
2 years ago
composer.json
2 years ago
FsProducer.php
110 lines
| 1 | <?php |
| 2 | |
| 3 | declare(strict_types=1); |
| 4 | |
| 5 | namespace Enqueue\Fs; |
| 6 | |
| 7 | use Interop\Queue\Destination; |
| 8 | use Interop\Queue\Exception\DeliveryDelayNotSupportedException; |
| 9 | use Interop\Queue\Exception\InvalidDestinationException; |
| 10 | use Interop\Queue\Exception\InvalidMessageException; |
| 11 | use Interop\Queue\Exception\PriorityNotSupportedException; |
| 12 | use Interop\Queue\Message; |
| 13 | use Interop\Queue\Producer; |
| 14 | use Makasim\File\TempFile; |
| 15 | |
| 16 | class FsProducer implements Producer |
| 17 | { |
| 18 | /** |
| 19 | * @var float|int|null |
| 20 | */ |
| 21 | private $timeToLive; |
| 22 | |
| 23 | /** |
| 24 | * @var FsContext |
| 25 | */ |
| 26 | private $context; |
| 27 | |
| 28 | public function __construct(FsContext $context) |
| 29 | { |
| 30 | $this->context = $context; |
| 31 | } |
| 32 | |
| 33 | /** |
| 34 | * @param FsDestination $destination |
| 35 | * @param FsMessage $message |
| 36 | */ |
| 37 | public function send(Destination $destination, Message $message): void |
| 38 | { |
| 39 | InvalidDestinationException::assertDestinationInstanceOf($destination, FsDestination::class); |
| 40 | InvalidMessageException::assertMessageInstanceOf($message, FsMessage::class); |
| 41 | |
| 42 | $this->context->workWithFile($destination, 'a+', function (FsDestination $destination, $file) use ($message) { |
| 43 | $fileInfo = $destination->getFileInfo(); |
| 44 | if ($fileInfo instanceof TempFile && false == file_exists((string) $fileInfo)) { |
| 45 | return; |
| 46 | } |
| 47 | |
| 48 | if (null !== $this->timeToLive) { |
| 49 | $message->setHeader('x-expire-at', microtime(true) + ($this->timeToLive / 1000)); |
| 50 | } |
| 51 | |
| 52 | $rawMessage = json_encode($message); |
| 53 | $rawMessage = str_replace('|{', '\|\{', $rawMessage); |
| 54 | $rawMessage = '|'.$rawMessage; |
| 55 | |
| 56 | if (JSON_ERROR_NONE !== json_last_error()) { |
| 57 | throw new \InvalidArgumentException(sprintf( |
| 58 | 'Could not encode value into json. Error %s and message %s', |
| 59 | json_last_error(), |
| 60 | json_last_error_msg() |
| 61 | )); |
| 62 | } |
| 63 | |
| 64 | $rawMessage = str_repeat(' ', 64 - (strlen($rawMessage) % 64)).$rawMessage; |
| 65 | |
| 66 | fwrite($file, $rawMessage); |
| 67 | }); |
| 68 | } |
| 69 | |
| 70 | public function setDeliveryDelay(int $deliveryDelay = null): Producer |
| 71 | { |
| 72 | if (null === $deliveryDelay) { |
| 73 | return $this; |
| 74 | } |
| 75 | |
| 76 | throw DeliveryDelayNotSupportedException::providerDoestNotSupportIt(); |
| 77 | } |
| 78 | |
| 79 | public function getDeliveryDelay(): ?int |
| 80 | { |
| 81 | return null; |
| 82 | } |
| 83 | |
| 84 | public function setPriority(int $priority = null): Producer |
| 85 | { |
| 86 | if (null === $priority) { |
| 87 | return $this; |
| 88 | } |
| 89 | |
| 90 | throw PriorityNotSupportedException::providerDoestNotSupportIt(); |
| 91 | } |
| 92 | |
| 93 | public function getPriority(): ?int |
| 94 | { |
| 95 | return null; |
| 96 | } |
| 97 | |
| 98 | public function setTimeToLive(int $timeToLive = null): Producer |
| 99 | { |
| 100 | $this->timeToLive = $timeToLive; |
| 101 | |
| 102 | return $this; |
| 103 | } |
| 104 | |
| 105 | public function getTimeToLive(): ?int |
| 106 | { |
| 107 | return null; |
| 108 | } |
| 109 | } |
| 110 |