]> git.mxchange.org Git - quix0rs-gnu-social.git/commitdiff
Several fixes to make RabbitMQ a player.
authorMarcel van der Boom <marcel@hsdev.com>
Tue, 8 Sep 2009 20:21:33 +0000 (22:21 +0200)
committerCraig Andrews <candrews@integralblue.com>
Sat, 12 Sep 2009 00:50:53 +0000 (20:50 -0400)
  * extlib/Stomp.php
    -spaces for tabs (we're on PEAR, right?)
    - send: initialize the $properties parameter as array() instead of null
      this prevents unsetting $headers if $properties was not set
      (besides that, it's the proper way to initialize an array)
    - subscribe: insert FIXME's on ActiveMQ specifics
    - ack: make sure the content-length header is set *and* is zero.
      I have seen the header set to '3' there but could not find where it
      came from, this is at least safe.
    - disconnect: typo in $headers variable
    - readFrame: use fgets() instead of gets() so that RabbitQ, which is more protocol strict can also play
  * extlib/Stomp/Frame.php
    - spaces for tabs
    - add note on possibly protocol violating linefeed
  * extlib/Stomp/Message.php
    - space for tabs
    - add content-length header for message
  * lib/stompqueuemanager.php
    - use the notice for logging, not the frame

extlib/Stomp.php
extlib/Stomp/Frame.php
extlib/Stomp/Message.php
lib/stompqueuemanager.php

index abd9cba62b44f633f97b75a2fce83424eb958748..c9e90629c4bde2f649e5422184e0dfe57adac1cf 100644 (file)
@@ -26,7 +26,7 @@ require_once 'Stomp/Frame.php';
  *
  * @package Stomp
  * @author Hiram Chirino <hiram@hiramchirino.com>
- * @author Dejan Bosanac <dejan@nighttale.net> 
+ * @author Dejan Bosanac <dejan@nighttale.net>
  * @author Michael Caplan <mcaplan@labnet.net>
  * @version $Revision: 43 $
  */
@@ -44,15 +44,15 @@ class Stomp
      *
      * @var int
      */
-       public $prefetchSize = 1;
-    
-       /**
+    public $prefetchSize = 1;
+
+    /**
      * Client id used for durable subscriptions
      *
      * @var string
      */
-       public $clientId = null;
-    
+    public $clientId = null;
+
     protected $_brokerUri = null;
     protected $_socket = null;
     protected $_hosts = array();
@@ -66,7 +66,7 @@ class Stomp
     protected $_sessionId;
     protected $_read_timeout_seconds = 60;
     protected $_read_timeout_milliseconds = 0;
-    
+
     /**
      * Constructor
      *
@@ -134,10 +134,10 @@ class Stomp
             require_once 'Stomp/Exception.php';
             throw new Stomp_Exception("No broker defined");
         }
-        
+
         // force disconnect, if previous established connection exists
         $this->disconnect();
-        
+
         $i = $this->_currentHost;
         $att = 0;
         $connected = false;
@@ -190,11 +190,11 @@ class Stomp
         if ($password != '') {
             $this->_password = $password;
         }
-               $headers = array('login' => $this->_username , 'passcode' => $this->_password);
-               if ($this->clientId != null) {
-                       $headers["client-id"] = $this->clientId;
-               }
-               $frame = new Stomp_Frame("CONNECT", $headers);
+        $headers = array('login' => $this->_username , 'passcode' => $this->_password);
+        if ($this->clientId != null) {
+            $headers["client-id"] = $this->clientId;
+        }
+        $frame = new Stomp_Frame("CONNECT", $headers);
         $this->_writeFrame($frame);
         $frame = $this->readFrame();
         if ($frame instanceof Stomp_Frame && $frame->command == 'CONNECTED') {
@@ -209,7 +209,7 @@ class Stomp
             }
         }
     }
-    
+
     /**
      * Check if client session has ben established
      *
@@ -229,7 +229,7 @@ class Stomp
         return $this->_sessionId;
     }
     /**
-     * Send a message to a destination in the messaging system 
+     * Send a message to a destination in the messaging system
      *
      * @param string $destination Destination queue
      * @param string|Stomp_Frame $msg Message
@@ -237,7 +237,7 @@ class Stomp
      * @param boolean $sync Perform request synchronously
      * @return boolean
      */
-    public function send ($destination, $msg, $properties = null, $sync = null)
+    public function send ($destination, $msg, $properties = array(), $sync = null)
     {
         if ($msg instanceof Stomp_Frame) {
             $msg->headers['destination'] = $destination;
@@ -319,10 +319,12 @@ class Stomp
     public function subscribe ($destination, $properties = null, $sync = null)
     {
         $headers = array('ack' => 'client');
-               $headers['activemq.prefetchSize'] = $this->prefetchSize;
-               if ($this->clientId != null) {
-                       $headers["activemq.subcriptionName"] = $this->clientId;
-               }
+        // FIXME: this seems to be activemq specific, but not hurting rabbitmq?
+        $headers['activemq.prefetchSize'] = $this->prefetchSize;
+        if ($this->clientId != null) {
+            // FIXME: this seems to be activemq specific, but not hurting rabbitmq?
+            $headers["activemq.subcriptionName"] = $this->clientId;
+        }
         if (isset($properties)) {
             foreach ($properties as $name => $value) {
                 $headers[$name] = $value;
@@ -424,7 +426,7 @@ class Stomp
     }
     /**
      * Acknowledge consumption of a message from a subscription
-        * Note: This operation is always asynchronous
+     * Note: This operation is always asynchronous
      *
      * @param string|Stomp_Frame $messageMessage ID
      * @param string $transactionId
@@ -433,20 +435,26 @@ class Stomp
      */
     public function ack ($message, $transactionId = null)
     {
+        // Handle the headers,
+        $headers = array();
+
         if ($message instanceof Stomp_Frame) {
-            $frame = new Stomp_Frame('ACK', $message->headers);
-            $this->_writeFrame($frame);
-            return true;
+            // Copy headers from the object
+            // FIXME: at least content-length can be wrong here (set to 3 sometimes).
+            $headers = $message->headers;
         } else {
-            $headers = array();
             if (isset($transactionId)) {
                 $headers['transaction'] = $transactionId;
             }
             $headers['message-id'] = $message;
-            $frame = new Stomp_Frame('ACK', $headers);
-            $this->_writeFrame($frame);
-            return true;
         }
+        // An ACK has no content
+        $headers['content-length'] = 0;
+
+        // Create it and write it out
+        $frame = new Stomp_Frame('ACK', $headers);
+        $this->_writeFrame($frame);
+        return true;
     }
     /**
      * Graceful disconnect from the server
@@ -454,11 +462,11 @@ class Stomp
      */
     public function disconnect ()
     {
-               $headers = array();
+        $headers = array();
 
-               if ($this->clientId != null) {
-                       $headers["client-id"] = $this->clientId;
-               }
+        if ($this->clientId != null) {
+            $headers["client-id"] = $this->clientId;
+        }
 
         if (is_resource($this->_socket)) {
             $this->_writeFrame(new Stomp_Frame('DISCONNECT', $headers));
@@ -490,19 +498,19 @@ class Stomp
             $this->_writeFrame($stompFrame);
         }
     }
-    
+
     /**
      * Set timeout to wait for content to read
      *
      * @param int $seconds_to_wait  Seconds to wait for a frame
      * @param int $milliseconds Milliseconds to wait for a frame
      */
-    public function setReadTimeout($seconds, $milliseconds = 0) 
+    public function setReadTimeout($seconds, $milliseconds = 0)
     {
         $this->_read_timeout_seconds = $seconds;
         $this->_read_timeout_milliseconds = $milliseconds;
     }
-    
+
     /**
      * Read responce frame from server
      *
@@ -513,19 +521,29 @@ class Stomp
         if (!$this->hasFrameToRead()) {
             return false;
         }
-        
+
         $rb = 1024;
         $data = '';
-        do {
-            $read = fgets($this->_socket, $rb);
-            if ($read === false) {
-                $this->_reconnect();
-                return $this->readFrame();
-            }
-            $data .= $read;
-            $len = strlen($data);
-        } while (($len < 2 || ! ($data[$len - 2] == "\x00" && $data[$len - 1] == "\n")));
-        
+         do {
+             $read = fread($this->_socket, $rb);
+             if ($read === false) {
+                 $this->_reconnect();
+                 return $this->readFrame();
+             }
+             $data .= $read;
+             $len = strlen($data);
+
+             $continue = true;
+             // ActiveMq apparently add \n after 0 char
+             if($data[$len - 2] == "\x00" && $data[$len - 1] == "\n")  {
+               $continue = false;
+             }
+
+             // RabbitMq does not
+             if($data[$len - 1] == "\x00") {
+                $continue = false;
+             }
+        } while ( $continue );
         list ($header, $body) = explode("\n\n", $data, 2);
         $header = explode("\n", $header);
         $headers = array();
@@ -546,7 +564,7 @@ class Stomp
             return $frame;
         }
     }
-    
+
     /**
      * Check if there is a frame to read
      *
@@ -557,7 +575,7 @@ class Stomp
         $read = array($this->_socket);
         $write = null;
         $except = null;
-        
+
         $has_frame_to_read = stream_select($read, $write, $except, $this->_read_timeout_seconds, $this->_read_timeout_milliseconds);
 
         if ($has_frame_to_read === false) {
@@ -565,18 +583,18 @@ class Stomp
         } else if ($has_frame_to_read > 0) {
             return true;
         } else {
-            return false; 
+            return false;
         }
     }
-    
+
     /**
      * Reconnects and renews subscriptions (if there were any)
-     * Call this method when you detect connection problems     
+     * Call this method when you detect connection problems
      */
     protected function _reconnect ()
     {
         $subscriptions = $this->_subscriptions;
-        
+
         $this->connect($this->_username, $this->_password);
         foreach ($subscriptions as $dest => $properties) {
             $this->subscribe($dest, $properties);
index dc59c1cb7fff04d6906496b39e8395f1e6a445c0..9fd97b4f5f766cd85b1db652f6286026e5d22b2c 100644 (file)
  */
 
 /* vim: set expandtab tabstop=3 shiftwidth=3: */
-\r
-/**\r
- * Stomp Frames are messages that are sent and received on a StompConnection.\r
- *\r
- * @package Stomp\r
- * @author Hiram Chirino <hiram@hiramchirino.com>\r
- * @author Dejan Bosanac <dejan@nighttale.net>\r
- * @author Michael Caplan <mcaplan@labnet.net>\r
- * @version $Revision: 36 $\r
- */\r
-class Stomp_Frame\r
-{\r
-    public $command;\r
-    public $headers = array();\r
-    public $body;\r
-    \r
-    /**\r
-     * Constructor\r
-     *\r
-     * @param string $command\r
-     * @param array $headers\r
-     * @param string $body\r
-     */\r
-    public function __construct ($command = null, $headers = null, $body = null)\r
-    {\r
-        $this->_init($command, $headers, $body);\r
-    }\r
-    \r
-    protected function _init ($command = null, $headers = null, $body = null)\r
-    {\r
-        $this->command = $command;\r
-        if ($headers != null) {\r
-            $this->headers = $headers;\r
-        }\r
-        $this->body = $body;\r
-        \r
-        if ($this->command == 'ERROR') {\r
-            require_once 'Stomp/Exception.php';\r
-            throw new Stomp_Exception($this->headers['message'], 0, $this->body);\r
-        }\r
+
+/**
+ * Stomp Frames are messages that are sent and received on a StompConnection.
+ *
+ * @package Stomp
+ * @author Hiram Chirino <hiram@hiramchirino.com>
+ * @author Dejan Bosanac <dejan@nighttale.net>
+ * @author Michael Caplan <mcaplan@labnet.net>
+ * @version $Revision: 36 $
+ */
+class Stomp_Frame
+{
+    public $command;
+    public $headers = array();
+    public $body;
+    
+    /**
+     * Constructor
+     *
+     * @param string $command
+     * @param array $headers
+     * @param string $body
+     */
+    public function __construct ($command = null, $headers = null, $body = null)
+    {
+        $this->_init($command, $headers, $body);
+    }
+    
+    protected function _init ($command = null, $headers = null, $body = null)
+    {
+        $this->command = $command;
+        if ($headers != null) {
+            $this->headers = $headers;
+        }
+        $this->body = $body;
+        
+        if ($this->command == 'ERROR') {
+            require_once 'Stomp/Exception.php';
+            throw new Stomp_Exception($this->headers['message'], 0, $this->body);
+        }
     }
     
     /**
@@ -74,7 +74,8 @@ class Stomp_Frame
         
         $data .= "\n";
         $data .= $this->body;
-        return $data .= "\x00\n";
-    }\r
-}\r
+        $data .= "\x00\n"; // Should there really be a linefeed here?
+        return $data;
+    }
+}
 ?>
\ No newline at end of file
index 6bcad3efd9c8525661715349fb27b05dc4559300..0556621338ad58b6993954ed8eca4896af955383 100644 (file)
@@ -29,8 +29,12 @@ require_once 'Stomp/Frame.php';
  */
 class Stomp_Message extends Stomp_Frame
 {
-    public function __construct ($body, $headers = null)
+    public function __construct ($body, $headers = array())
     {
+        if(!isset($headers['content-length'])) {
+          // TODO: log this, to see if this is correct
+          $headers['content-length'] = strlen($body);
+        }
         $this->_init("SEND", $headers, $body);
     }
 }
index f059b42f0095f69f5db827bbf3f5698d128f87ba..5d8b2996b272b2ecc6f4217ef55e47a36e20b126 100644 (file)
@@ -141,10 +141,11 @@ class StompQueueManager
                 $this->con->ack($frame);
             } else {
                 if ($handler->handle_notice($notice)) {
-                    $this->_log(LOG_INFO, 'Successfully handled notice '. $notice->id .' posted at ' . $frame->headers['created'] . ' in queue '. $queue);
+                    $this->_log(LOG_INFO, 'Successfully handled notice '. $notice->id .' originally posted at ' . $notice->created . ' in queue '. $queue);
+                    
                     $this->con->ack($frame);
                 } else {
-                    $this->_log(LOG_WARNING, 'Failed handling notice '. $notice->id .' posted at ' . $frame->headers['created']  . ' in queue '. $queue);
+                    $this->_log(LOG_WARNING, 'Failed handling notice '. $notice->id .' originally posted at ' . $notice->created   . ' in queue '. $queue);
                     // FIXME we probably shouldn't have to do
                     // this kind of queue management ourselves
                     $this->con->ack($frame);