|
1
|
|
|
<?php |
|
2
|
|
|
|
|
3
|
|
|
namespace Basis; |
|
4
|
|
|
|
|
5
|
|
|
use Exception; |
|
6
|
|
|
use Tarantool\Client\Client; |
|
7
|
|
|
use Tarantool\Client\Connection\StreamConnection; |
|
8
|
|
|
use Tarantool\Client\Packer\PurePacker; |
|
9
|
|
|
|
|
10
|
|
|
class Queue extends Client |
|
11
|
|
|
{ |
|
12
|
1 |
|
public function __construct(Config $config) |
|
13
|
|
|
{ |
|
14
|
1 |
|
$connection = new StreamConnection('tcp://'.$config['queue.host'].':'.$config['queue.port'], [ |
|
15
|
1 |
|
'socket_timeout' => $config['queue.socket_timeout'] ?: 5, |
|
16
|
1 |
|
'connect_timeout' => $config['queue.connect_timeout'] ?: 5, |
|
17
|
|
|
]); |
|
18
|
1 |
|
parent::__construct($connection, new PurePacker()); |
|
19
|
|
|
} |
|
20
|
1 |
|
|
|
21
|
1 |
|
public function init($tube, $type = 'fifottl') |
|
22
|
|
|
{ |
|
23
|
1 |
|
return $this->evaluate(" |
|
24
|
|
|
box.once('$tube-tube', function() |
|
25
|
1 |
|
local queue = require('queue') |
|
26
|
|
|
queue.create_tube('$tube', '$type') |
|
27
|
|
|
end) |
|
28
|
1 |
|
"); |
|
29
|
|
|
} |
|
30
|
1 |
|
|
|
31
|
1 |
|
public function truncate($tube) |
|
32
|
|
|
{ |
|
33
|
1 |
|
$this->evaluate("require('queue').tube.$tube:truncate()"); |
|
34
|
|
|
} |
|
35
|
1 |
|
|
|
36
|
1 |
|
public function take($tube, $timeout = 1) |
|
37
|
1 |
|
{ |
|
38
|
|
|
$tasks = $this->evaluate("return require('queue').tube.$tube:take($timeout)")->getData(); |
|
39
|
1 |
|
if(count($tasks)) { |
|
40
|
|
|
return $tasks[0]; |
|
41
|
1 |
|
} |
|
42
|
|
|
} |
|
43
|
1 |
|
|
|
44
|
|
|
public function ack($tube, $task) |
|
45
|
|
|
{ |
|
46
|
1 |
|
return $this->evaluate("require('queue').tube.$tube:ack($task)"); |
|
47
|
|
|
} |
|
48
|
1 |
|
|
|
49
|
|
|
public function bury($tube, $task) |
|
50
|
|
|
{ |
|
51
|
|
|
return $this->evaluate("require('queue').tube.$tube:bury($task)"); |
|
52
|
|
|
} |
|
53
|
|
|
|
|
54
|
|
|
public function put($tube, $task, $options = []) |
|
55
|
|
|
{ |
|
56
|
|
|
return $this->evaluate("require('queue').tube.$tube:put(...)", [$task, $options]); |
|
57
|
|
|
} |
|
58
|
|
|
|
|
59
|
|
|
public function release($tube, $id, $options = []) |
|
60
|
|
|
{ |
|
61
|
|
|
return $this->evaluate("require('queue').tube.$tube:release($id, ...)", [$options]); |
|
62
|
|
|
} |
|
63
|
|
|
} |