.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
FsConnectionFactory.php
102 lines
| 1 | <?php |
| 2 | |
| 3 | declare(strict_types=1); |
| 4 | |
| 5 | namespace Enqueue\Fs; |
| 6 | |
| 7 | use Enqueue\Dsn\Dsn; |
| 8 | use Interop\Queue\ConnectionFactory; |
| 9 | use Interop\Queue\Context; |
| 10 | |
| 11 | class FsConnectionFactory implements ConnectionFactory |
| 12 | { |
| 13 | /** |
| 14 | * @var string |
| 15 | */ |
| 16 | private $config; |
| 17 | |
| 18 | /** |
| 19 | * The config could be an array, string DSN or null. In case of null it will attempt to store files in /tmp/enqueue folder. |
| 20 | * |
| 21 | * [ |
| 22 | * 'path' => 'the directory where all queue\topic files remain. For example /home/foo/enqueue', |
| 23 | * 'pre_fetch_count' => 'Integer. Defines how many messages to fetch from the file.', |
| 24 | * 'chmod' => 'Defines a mode the files are created with', |
| 25 | * 'polling_interval' => 'How often query for new messages, default 100 (milliseconds)', |
| 26 | * ] |
| 27 | * |
| 28 | * or |
| 29 | * |
| 30 | * file: - create queue files in tmp dir. |
| 31 | * file:///home/foo/enqueue |
| 32 | * file:///home/foo/enqueue?pre_fetch_count=20&chmod=0777 |
| 33 | * |
| 34 | * @param array|string|null $config |
| 35 | */ |
| 36 | public function __construct($config = 'file:') |
| 37 | { |
| 38 | if (empty($config) || 'file:' === $config) { |
| 39 | $config = $this->parseDsn('file://'.sys_get_temp_dir().'/enqueue'); |
| 40 | } elseif (is_string($config)) { |
| 41 | if ('/' === $config[0]) { |
| 42 | $config = $this->parseDsn('file://'.$config); |
| 43 | } else { |
| 44 | $config = $this->parseDsn($config); |
| 45 | } |
| 46 | } elseif (is_array($config)) { |
| 47 | } else { |
| 48 | throw new \LogicException('The config must be either an array of options, a DSN string or null'); |
| 49 | } |
| 50 | |
| 51 | $this->config = array_replace($this->defaultConfig(), $config); |
| 52 | |
| 53 | if (empty($this->config['path'])) { |
| 54 | throw new \LogicException('The path option must be set.'); |
| 55 | } |
| 56 | } |
| 57 | |
| 58 | /** |
| 59 | * @return FsContext |
| 60 | */ |
| 61 | public function createContext(): Context |
| 62 | { |
| 63 | return new FsContext( |
| 64 | $this->config['path'], |
| 65 | $this->config['pre_fetch_count'], |
| 66 | $this->config['chmod'], |
| 67 | $this->config['polling_interval'] |
| 68 | ); |
| 69 | } |
| 70 | |
| 71 | private function parseDsn(string $dsn): array |
| 72 | { |
| 73 | $dsn = Dsn::parseFirst($dsn); |
| 74 | |
| 75 | $supportedSchemes = ['file']; |
| 76 | if (false == in_array($dsn->getSchemeProtocol(), $supportedSchemes, true)) { |
| 77 | throw new \LogicException(sprintf( |
| 78 | 'The given scheme protocol "%s" is not supported. It must be one of "%s"', |
| 79 | $dsn->getSchemeProtocol(), |
| 80 | implode('", "', $supportedSchemes) |
| 81 | )); |
| 82 | } |
| 83 | |
| 84 | return array_filter(array_replace($dsn->getQuery(), [ |
| 85 | 'path' => $dsn->getPath(), |
| 86 | 'pre_fetch_count' => $dsn->getDecimal('pre_fetch_count'), |
| 87 | 'chmod' => $dsn->getOctal('chmod'), |
| 88 | 'polling_interval' => $dsn->getDecimal('polling_interval'), |
| 89 | ]), function ($value) { return null !== $value; }); |
| 90 | } |
| 91 | |
| 92 | private function defaultConfig(): array |
| 93 | { |
| 94 | return [ |
| 95 | 'path' => null, |
| 96 | 'pre_fetch_count' => 1, |
| 97 | 'chmod' => 0600, |
| 98 | 'polling_interval' => 100, |
| 99 | ]; |
| 100 | } |
| 101 | } |
| 102 |