1
|
|
|
<?php |
2
|
|
|
|
3
|
|
|
namespace Bernard\Driver\NewMongoDB; |
4
|
|
|
|
5
|
|
|
use MongoDB\Collection; |
6
|
|
|
use MongoDB\BSON\UTCDateTime; |
7
|
|
|
use MongoDB\BSON\ObjectID; |
8
|
|
|
|
9
|
|
|
/** |
10
|
|
|
* Driver supporting MongoDB. |
11
|
|
|
*/ |
12
|
|
|
final class Driver implements \Bernard\Driver |
13
|
|
|
{ |
14
|
|
|
private $messages; |
15
|
|
|
private $queues; |
16
|
|
|
|
17
|
|
|
/** |
18
|
|
|
* @param MongoCollection $queues Collection where queues will be stored |
19
|
|
|
* @param MongoCollection $messages Collection where messages will be stored |
20
|
|
|
*/ |
21
|
|
|
public function __construct(Collection $queues, Collection $messages) |
22
|
|
|
{ |
23
|
|
|
$this->queues = $queues; |
24
|
|
|
$this->messages = $messages; |
25
|
|
|
} |
26
|
|
|
|
27
|
|
|
/** |
28
|
|
|
* {@inheritdoc} |
29
|
|
|
*/ |
30
|
|
|
public function listQueues() |
31
|
|
|
{ |
32
|
|
|
return $this->queues->distinct('_id'); |
33
|
|
|
} |
34
|
|
|
|
35
|
|
|
/** |
36
|
|
|
* {@inheritdoc} |
37
|
|
|
*/ |
38
|
|
|
public function createQueue($queueName) |
39
|
|
|
{ |
40
|
|
|
$data = ['_id' => (string) $queueName]; |
41
|
|
|
$updateData = ['$set' => $data]; |
42
|
|
|
|
43
|
|
|
$this->queues->updateOne($data, $updateData, ['upsert' => true]); |
44
|
|
|
} |
45
|
|
|
|
46
|
|
|
/** |
47
|
|
|
* {@inheritdoc} |
48
|
|
|
*/ |
49
|
|
|
public function countMessages($queueName) |
50
|
|
|
{ |
51
|
|
|
return $this->messages->count([ |
52
|
|
|
'queue' => (string) $queueName, |
53
|
|
|
'visible' => true, |
54
|
|
|
]); |
55
|
|
|
} |
56
|
|
|
|
57
|
|
|
/** |
58
|
|
|
* {@inheritdoc} |
59
|
|
|
*/ |
60
|
|
View Code Duplication |
public function pushMessage($queueName, $message) |
|
|
|
|
61
|
|
|
{ |
62
|
|
|
$data = [ |
63
|
|
|
'queue' => (string) $queueName, |
64
|
|
|
'message' => (string) $message, |
65
|
|
|
'sentAt' => new UTCDateTime(), |
66
|
|
|
'visible' => true, |
67
|
|
|
]; |
68
|
|
|
|
69
|
|
|
$this->messages->insertOne($data); |
70
|
|
|
} |
71
|
|
|
|
72
|
|
|
/** |
73
|
|
|
* {@inheritdoc} |
74
|
|
|
*/ |
75
|
|
View Code Duplication |
public function popMessage($queueName, $duration = 5) |
|
|
|
|
76
|
|
|
{ |
77
|
|
|
$runtime = microtime(true) + $duration; |
78
|
|
|
|
79
|
|
|
while (microtime(true) < $runtime) { |
80
|
|
|
$result = $this->messages->findOneAndUpdate( |
81
|
|
|
['queue' => (string) $queueName, 'visible' => true], |
82
|
|
|
['$set' => ['visible' => false]], |
83
|
|
|
['sort' => ['sentAt' => 1], 'projection' => ['message' => 1]] |
84
|
|
|
); |
85
|
|
|
|
86
|
|
|
if ($result) { |
87
|
|
|
return [(string) $result['message'], (string) $result['_id']]; |
88
|
|
|
} |
89
|
|
|
|
90
|
|
|
usleep(10000); |
91
|
|
|
} |
92
|
|
|
|
93
|
|
|
return [null, null]; |
94
|
|
|
} |
95
|
|
|
|
96
|
|
|
/** |
97
|
|
|
* {@inheritdoc} |
98
|
|
|
*/ |
99
|
|
|
public function acknowledgeMessage($queueName, $receipt) |
100
|
|
|
{ |
101
|
|
|
$this->messages->deleteOne([ |
102
|
|
|
'_id' => new ObjectId((string) $receipt), |
103
|
|
|
'queue' => (string) $queueName, |
104
|
|
|
]); |
105
|
|
|
} |
106
|
|
|
|
107
|
|
|
/** |
108
|
|
|
* {@inheritdoc} |
109
|
|
|
*/ |
110
|
|
|
public function peekQueue($queueName, $index = 0, $limit = 20) |
111
|
|
|
{ |
112
|
|
|
$query = ['queue' => (string) $queueName, 'visible' => true]; |
113
|
|
|
$fields = ['_id' => 0, 'message' => 1]; |
114
|
|
|
|
115
|
|
|
$results = $this->messages |
116
|
|
|
->find( |
117
|
|
|
$query, |
118
|
|
|
[ |
119
|
|
|
'projection' => $fields, |
120
|
|
|
'sort' => ['sentAt' => 1], |
121
|
|
|
'limit' => $limit, |
122
|
|
|
'skip' => $index |
123
|
|
|
] |
124
|
|
|
) |
125
|
|
|
->toArray(); |
126
|
|
|
|
127
|
|
|
return array_map(function($result){return $result['message']; }, $results); |
128
|
|
|
} |
129
|
|
|
|
130
|
|
|
/** |
131
|
|
|
* {@inheritdoc} |
132
|
|
|
*/ |
133
|
|
|
public function removeQueue($queueName) |
134
|
|
|
{ |
135
|
|
|
$this->queues->deleteOne(['_id' => $queueName]); |
136
|
|
|
$this->messages->deleteMany(['queue' => (string) $queueName]); |
137
|
|
|
} |
138
|
|
|
|
139
|
|
|
/** |
140
|
|
|
* {@inheritdoc} |
141
|
|
|
*/ |
142
|
|
|
public function info() |
143
|
|
|
{ |
144
|
|
|
return [ |
145
|
|
|
'messages' => (string) $this->messages, |
146
|
|
|
'queues' => (string) $this->queues, |
147
|
|
|
]; |
148
|
|
|
} |
149
|
|
|
} |
150
|
|
|
|
Duplicated code is one of the most pungent code smells. If you need to duplicate the same code in three or more different places, we strongly encourage you to look into extracting the code into a single class or operation.
You can also find more detailed suggestions in the “Code” section of your repository.