-
Notifications
You must be signed in to change notification settings - Fork 129
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Merge pull request #133 from jmikola/mongodb
MongoDB storage driver
- Loading branch information
Showing
5 changed files
with
541 additions
and
0 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,150 @@ | ||
<?php | ||
|
||
namespace Bernard\Driver; | ||
|
||
use MongoCollection; | ||
use MongoDate; | ||
use MongoId; | ||
|
||
/** | ||
* Driver supporting MongoDB | ||
* | ||
* @package Bernard | ||
*/ | ||
class MongoDBDriver implements \Bernard\Driver | ||
{ | ||
private $messages; | ||
private $queues; | ||
|
||
/** | ||
* Constructor. | ||
* | ||
* @param MongoCollection $queues Collection where queues will be stored | ||
* @param MongoCollection $messages Collection where messages will be stored | ||
*/ | ||
public function __construct(MongoCollection $queues, MongoCollection $messages) | ||
{ | ||
$this->queues = $queues; | ||
$this->messages = $messages; | ||
} | ||
|
||
/** | ||
* {@inheritDoc} | ||
*/ | ||
public function listQueues() | ||
{ | ||
return $this->queues->distinct('_id'); | ||
} | ||
|
||
/** | ||
* {@inheritDoc} | ||
*/ | ||
public function createQueue($queueName) | ||
{ | ||
$data = array('_id' => (string) $queueName); | ||
|
||
$this->queues->update($data, $data, array('upsert' => true)); | ||
} | ||
|
||
/** | ||
* {@inheritDoc} | ||
*/ | ||
public function countMessages($queueName) | ||
{ | ||
return $this->messages->count(array( | ||
'queue' => (string) $queueName, | ||
'visible' => true, | ||
)); | ||
} | ||
|
||
/** | ||
* {@inheritDoc} | ||
*/ | ||
public function pushMessage($queueName, $message) | ||
{ | ||
$data = array( | ||
'queue' => (string) $queueName, | ||
'message' => (string) $message, | ||
'sentAt' => new MongoDate(), | ||
'visible' => true, | ||
); | ||
|
||
$this->messages->insert($data); | ||
} | ||
|
||
/** | ||
* {@inheritDoc} | ||
*/ | ||
public function popMessage($queueName, $interval = 5) | ||
{ | ||
$runtime = microtime(true) + $interval; | ||
|
||
while (microtime(true) < $runtime) { | ||
$result = $this->messages->findAndModify( | ||
array('queue' => (string) $queueName, 'visible' => true), | ||
array('$set' => array('visible' => false)), | ||
array('message' => 1), | ||
array('sort' => array('sentAt' => 1)) | ||
); | ||
|
||
if ($result) { | ||
return array((string) $result['message'], (string) $result['_id']); | ||
} | ||
|
||
usleep(10000); | ||
} | ||
|
||
return array(null, null); | ||
} | ||
|
||
/** | ||
* {@inheritDoc} | ||
*/ | ||
public function acknowledgeMessage($queueName, $receipt) | ||
{ | ||
$this->messages->remove(array( | ||
'_id' => new MongoId((string) $receipt), | ||
'queue' => (string) $queueName, | ||
)); | ||
} | ||
|
||
/** | ||
* {@inheritDoc} | ||
*/ | ||
public function peekQueue($queueName, $index = 0, $limit = 20) | ||
{ | ||
$cursor = $this->messages->find( | ||
array('queue' => (string) $queueName, 'visible' => true), | ||
array('_id' => 0, 'message' => 1) | ||
) | ||
->sort(array('sentAt' => 1)) | ||
->limit($limit) | ||
->skip($index) | ||
; | ||
|
||
return array_map( | ||
function ($result) { return (string) $result['message']; }, | ||
iterator_to_array($cursor, false) | ||
); | ||
} | ||
|
||
/** | ||
* {@inheritDoc} | ||
*/ | ||
public function removeQueue($queueName) | ||
{ | ||
$this->queues->remove(array('_id' => $queueName)); | ||
$this->messages->remove(array('queue' => (string) $queueName)); | ||
} | ||
|
||
/** | ||
* {@inheritDoc} | ||
*/ | ||
public function info() | ||
{ | ||
return array( | ||
'messages' => (string) $this->messages, | ||
'queues' => (string) $this->queues, | ||
); | ||
} | ||
} |
154 changes: 154 additions & 0 deletions
154
tests/Bernard/Tests/Driver/MongoDBDriverFunctionalTest.php
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,154 @@ | ||
<?php | ||
|
||
namespace Bernard\Tests\Driver; | ||
|
||
use Bernard\Driver\MongoDBDriver; | ||
use MongoClient; | ||
use MongoCollection; | ||
use MongoException; | ||
|
||
/** | ||
* @coversDefaultClass Bernard\Driver\MongoDBDriver | ||
*/ | ||
class MongoDBDriverFunctionalTest extends \PHPUnit_Framework_TestCase | ||
{ | ||
const DATABASE = 'bernardQueueTest'; | ||
const MESSAGES = 'bernardMessages'; | ||
const QUEUES = 'bernardQueues'; | ||
|
||
private $messages; | ||
private $queues; | ||
private $driver; | ||
|
||
public function setUp() | ||
{ | ||
if ( ! class_exists('MongoClient')) { | ||
$this->markTestSkipped('MongoDB extension is not available.'); | ||
} | ||
|
||
try { | ||
$mongoClient = new MongoClient(); | ||
} catch (MongoConnectionException $e) { | ||
$this->markTestSkipped('Cannot connect to MongoDB server.'); | ||
} | ||
|
||
$this->queues = $mongoClient->selectCollection(self::DATABASE, self::QUEUES); | ||
$this->messages = $mongoClient->selectCollection(self::DATABASE, self::MESSAGES); | ||
$this->driver = new MongoDBDriver($this->queues, $this->messages); | ||
} | ||
|
||
public function tearDown() | ||
{ | ||
if ( ! $this->messages instanceof MongoCollection) { | ||
return; | ||
} | ||
|
||
$this->messages->drop(); | ||
$this->queues->drop(); | ||
} | ||
|
||
/** | ||
* @medium | ||
* @covers ::acknowledgeMessage() | ||
* @covers ::countMessages() | ||
* @covers ::popMessage() | ||
* @covers ::pushMessage() | ||
*/ | ||
public function testMessageLifecycle() | ||
{ | ||
$this->assertEquals(0, $this->driver->countMessages('foo')); | ||
|
||
$this->driver->pushMessage('foo', 'message1'); | ||
$this->assertEquals(1, $this->driver->countMessages('foo')); | ||
|
||
$this->driver->pushMessage('foo', 'message2'); | ||
$this->assertEquals(2, $this->driver->countMessages('foo')); | ||
|
||
list($message1, $receipt1) = $this->driver->popMessage('foo'); | ||
$this->assertSame('message1', $message1, 'The first message pushed is popped first'); | ||
$this->assertRegExp('/^[a-f\d]{24}$/i', $receipt1, 'The message receipt is an ObjectId'); | ||
$this->assertEquals(1, $this->driver->countMessages('foo')); | ||
|
||
list($message2, $receipt2) = $this->driver->popMessage('foo'); | ||
$this->assertSame('message2', $message2, 'The second message pushed is popped second'); | ||
$this->assertRegExp('/^[a-f\d]{24}$/i', $receipt2, 'The message receipt is an ObjectId'); | ||
$this->assertEquals(0, $this->driver->countMessages('foo')); | ||
|
||
list($message3, $receipt3) = $this->driver->popMessage('foo', 1); | ||
$this->assertNull($message3, 'Null message is returned when popping an empty queue'); | ||
$this->assertNull($receipt3, 'Null receipt is returned when popping an empty queue'); | ||
|
||
$this->assertEquals(2, $this->messages->count(), 'Popped messages remain in the database'); | ||
|
||
$this->driver->acknowledgeMessage('foo', $receipt1); | ||
$this->assertEquals(1, $this->messages->count(), 'Acknowledged messages are removed from the database'); | ||
|
||
$this->driver->acknowledgeMessage('foo', $receipt2); | ||
$this->assertEquals(0, $this->messages->count(), 'Acknowledged messages are removed from the database'); | ||
} | ||
|
||
public function testPeekQueue() | ||
{ | ||
$this->driver->pushMessage('foo', 'message1'); | ||
$this->driver->pushMessage('foo', 'message2'); | ||
|
||
$this->assertSame(array('message1', 'message2'), $this->driver->peekQueue('foo')); | ||
$this->assertSame(array('message2'), $this->driver->peekQueue('foo', 1)); | ||
$this->assertSame(array(), $this->driver->peekQueue('foo', 2)); | ||
$this->assertSame(array('message1'), $this->driver->peekQueue('foo', 0, 1)); | ||
$this->assertSame(array('message2'), $this->driver->peekQueue('foo', 1, 1)); | ||
} | ||
|
||
/** | ||
* @covers ::createQueue() | ||
* @covers ::listQueues() | ||
* @covers ::removeQueue() | ||
*/ | ||
public function testQueueLifecycle() | ||
{ | ||
$this->driver->createQueue('foo'); | ||
$this->driver->createQueue('bar'); | ||
|
||
$queues = $this->driver->listQueues(); | ||
$this->assertCount(2, $queues); | ||
$this->assertContains('foo', $queues); | ||
$this->assertContains('bar', $queues); | ||
|
||
$this->driver->removeQueue('foo'); | ||
|
||
$queues = $this->driver->listQueues(); | ||
$this->assertCount(1, $queues); | ||
$this->assertNotContains('foo', $queues); | ||
$this->assertContains('bar', $queues); | ||
} | ||
|
||
public function testRemoveQueueDeletesMessages() | ||
{ | ||
$this->driver->pushMessage('foo', 'message1'); | ||
$this->driver->pushMessage('foo', 'message2'); | ||
$this->assertEquals(2, $this->driver->countMessages('foo')); | ||
$this->assertEquals(2, $this->messages->count()); | ||
|
||
$this->driver->removeQueue('foo'); | ||
$this->assertEquals(0, $this->driver->countMessages('foo')); | ||
$this->assertEquals(0, $this->messages->count()); | ||
} | ||
|
||
public function testCreateQueueWithDuplicateNameIsNoop() | ||
{ | ||
$this->driver->createQueue('foo'); | ||
$this->driver->createQueue('foo'); | ||
|
||
$this->assertSame(array('foo'), $this->driver->listQueues()); | ||
} | ||
|
||
public function testInfo() | ||
{ | ||
$info = array( | ||
'messages' => self::DATABASE . '.' . self::MESSAGES, | ||
'queues' => self::DATABASE . '.' . self::QUEUES, | ||
); | ||
|
||
$this->assertSame($info, $this->driver->info()); | ||
} | ||
} |
Oops, something went wrong.