| 1 |
<?php |
| 2 |
|
| 3 |
if (!function_exists('json_encode')) { |
| 4 |
throw new Exception('The JSON PHP extension is required.'); |
| 5 |
} |
| 6 |
|
| 7 |
/** |
| 8 |
* Provides some base methods for use by a message Producer |
| 9 |
*/ |
| 10 |
abstract class Imagify_WPMedia_Producers_MixpanelBaseProducer extends Imagify_WPMedia_Base_MixpanelBase { |
| 11 |
|
| 12 |
|
| 13 |
/** |
| 14 |
* @var string a token associated to a Mixpanel project |
| 15 |
*/ |
| 16 |
protected $_token; |
| 17 |
|
| 18 |
|
| 19 |
/** |
| 20 |
* @var array a queue to hold messages in memory before flushing in batches |
| 21 |
*/ |
| 22 |
private $_queue = array(); |
| 23 |
|
| 24 |
|
| 25 |
/** |
| 26 |
* @var Imagify_WPMedia_ConsumerStrategies_AbstractConsumer the consumer to use when flushing messages |
| 27 |
*/ |
| 28 |
private $_consumer = null; |
| 29 |
|
| 30 |
|
| 31 |
/** |
| 32 |
* @var array The list of available consumers |
| 33 |
*/ |
| 34 |
private $_consumers = array( |
| 35 |
"file" => "Imagify_WPMedia_ConsumerStrategies_FileConsumer", |
| 36 |
"curl" => "Imagify_WPMedia_ConsumerStrategies_CurlConsumer", |
| 37 |
"socket" => "Imagify_WPMedia_ConsumerStrategies_SocketConsumer", |
| 38 |
); |
| 39 |
|
| 40 |
|
| 41 |
/** |
| 42 |
* If the queue reaches this size we'll auto-flush to prevent out of memory errors |
| 43 |
* @var int |
| 44 |
*/ |
| 45 |
protected $_max_queue_size = 1000; |
| 46 |
|
| 47 |
|
| 48 |
/** |
| 49 |
* Creates a new MixpanelBaseProducer, assings Mixpanel project token, registers custom Consumers, and instantiates |
| 50 |
* the desired consumer |
| 51 |
* @param $token |
| 52 |
* @param array $options |
| 53 |
*/ |
| 54 |
public function __construct($token, $options = array()) { |
| 55 |
|
| 56 |
parent::__construct($options); |
| 57 |
|
| 58 |
// register any customer consumers |
| 59 |
if (isset($options["consumers"])) { |
| 60 |
$this->_consumers = array_merge($this->_consumers, $options['consumers']); |
| 61 |
} |
| 62 |
|
| 63 |
// set max queue size |
| 64 |
if (isset($options["max_queue_size"])) { |
| 65 |
$this->_max_queue_size = $options['max_queue_size']; |
| 66 |
} |
| 67 |
|
| 68 |
// associate token |
| 69 |
$this->_token = $token; |
| 70 |
|
| 71 |
if ($this->_debug()) { |
| 72 |
$this->_log("Using token: ".$this->_token); |
| 73 |
} |
| 74 |
|
| 75 |
// instantiate the chosen consumer |
| 76 |
$this->_consumer = $this->_getConsumer(); |
| 77 |
|
| 78 |
} |
| 79 |
|
| 80 |
|
| 81 |
/** |
| 82 |
* Flush the queue when we destruct the client with retries |
| 83 |
*/ |
| 84 |
public function __destruct() { |
| 85 |
$attempts = 0; |
| 86 |
$max_attempts = 10; |
| 87 |
$success = false; |
| 88 |
while (!$success && $attempts < $max_attempts) { |
| 89 |
if ($this->_debug()) { |
| 90 |
$this->_log("destruct flush attempt #".($attempts+1)); |
| 91 |
} |
| 92 |
$success = $this->flush(); |
| 93 |
$attempts++; |
| 94 |
} |
| 95 |
} |
| 96 |
|
| 97 |
|
| 98 |
/** |
| 99 |
* Iterate the queue and write in batches using the instantiated Consumer Strategy |
| 100 |
* @param int $desired_batch_size |
| 101 |
* @return bool whether or not the flush was successful |
| 102 |
*/ |
| 103 |
public function flush($desired_batch_size = 50) { |
| 104 |
$queue_size = count($this->_queue); |
| 105 |
$succeeded = true; |
| 106 |
$num_threads = $this->_consumer->getNumThreads(); |
| 107 |
|
| 108 |
if ($this->_debug()) { |
| 109 |
$this->_log("Flush called - queue size: ".$queue_size); |
| 110 |
} |
| 111 |
|
| 112 |
while($queue_size > 0 && $succeeded) { |
| 113 |
$batch_size = min(array($queue_size, $desired_batch_size*$num_threads, $this->_options['max_batch_size']*$num_threads)); |
| 114 |
$batch = array_splice($this->_queue, 0, $batch_size); |
| 115 |
$succeeded = $this->_persist($batch); |
| 116 |
|
| 117 |
if (!$succeeded) { |
| 118 |
if ($this->_debug()) { |
| 119 |
$this->_log("Batch consumption failed!"); |
| 120 |
} |
| 121 |
$this->_queue = array_merge($batch, $this->_queue); |
| 122 |
|
| 123 |
if ($this->_debug()) { |
| 124 |
$this->_log("added batch back to queue, queue size is now $queue_size"); |
| 125 |
} |
| 126 |
} |
| 127 |
|
| 128 |
$queue_size = count($this->_queue); |
| 129 |
|
| 130 |
if ($this->_debug()) { |
| 131 |
$this->_log("Batch of $batch_size consumed, queue size is now $queue_size"); |
| 132 |
} |
| 133 |
} |
| 134 |
return $succeeded; |
| 135 |
} |
| 136 |
|
| 137 |
|
| 138 |
/** |
| 139 |
* Empties the queue without persisting any of the messages |
| 140 |
*/ |
| 141 |
public function reset() { |
| 142 |
$this->_queue = array(); |
| 143 |
} |
| 144 |
|
| 145 |
|
| 146 |
/** |
| 147 |
* Returns the in-memory queue |
| 148 |
* @return array |
| 149 |
*/ |
| 150 |
public function getQueue() { |
| 151 |
return $this->_queue; |
| 152 |
} |
| 153 |
|
| 154 |
/** |
| 155 |
* Returns the current Mixpanel project token |
| 156 |
* @return string |
| 157 |
*/ |
| 158 |
public function getToken() { |
| 159 |
return $this->_token; |
| 160 |
} |
| 161 |
|
| 162 |
|
| 163 |
/** |
| 164 |
* Given a strategy type, return a new PersistenceStrategy object |
| 165 |
* @return Imagify_WPMedia_ConsumerStrategies_AbstractConsumer |
| 166 |
*/ |
| 167 |
protected function _getConsumer() { |
| 168 |
$key = $this->_options['consumer']; |
| 169 |
$Strategy = $this->_consumers[$key]; |
| 170 |
if ($this->_debug()) { |
| 171 |
$this->_log("Using consumer: " . $key . " -> " . $Strategy); |
| 172 |
} |
| 173 |
$this->_options['endpoint'] = $this->_getEndpoint(); |
| 174 |
|
| 175 |
return new $Strategy($this->_options); |
| 176 |
} |
| 177 |
|
| 178 |
|
| 179 |
/** |
| 180 |
* Add an array representing a message to be sent to Mixpanel to a queue. |
| 181 |
* @param array $message |
| 182 |
*/ |
| 183 |
public function enqueue($message = array()) { |
| 184 |
array_push($this->_queue, $message); |
| 185 |
|
| 186 |
// force a flush if we've reached our threshold |
| 187 |
if (count($this->_queue) > $this->_max_queue_size) { |
| 188 |
$this->flush(); |
| 189 |
} |
| 190 |
|
| 191 |
if ($this->_debug()) { |
| 192 |
$this->_log("Queued message: ".json_encode($message)); |
| 193 |
} |
| 194 |
} |
| 195 |
|
| 196 |
|
| 197 |
/** |
| 198 |
* Add an array representing a list of messages to be sent to Mixpanel to a queue. |
| 199 |
* @param array $messages |
| 200 |
*/ |
| 201 |
public function enqueueAll($messages = array()) { |
| 202 |
foreach($messages as $message) { |
| 203 |
$this->enqueue($message); |
| 204 |
} |
| 205 |
|
| 206 |
} |
| 207 |
|
| 208 |
|
| 209 |
/** |
| 210 |
* Given an array of messages, persist it with the instantiated Persistence Strategy |
| 211 |
* @param $message |
| 212 |
* @return mixed |
| 213 |
*/ |
| 214 |
protected function _persist($message) { |
| 215 |
return $this->_consumer->persist($message); |
| 216 |
} |
| 217 |
|
| 218 |
|
| 219 |
|
| 220 |
|
| 221 |
/** |
| 222 |
* Return the endpoint that should be used by a consumer that consumes messages produced by this producer. |
| 223 |
* @return string |
| 224 |
*/ |
| 225 |
abstract function _getEndpoint(); |
| 226 |
|
| 227 |
} |
| 228 |
|