| 1 |
<?php |
| 2 |
|
| 3 |
/** |
| 4 |
* Consumes messages and sends them to a host/endpoint using cURL |
| 5 |
*/ |
| 6 |
class Imagify_WPMedia_ConsumerStrategies_CurlConsumer extends Imagify_WPMedia_ConsumerStrategies_AbstractConsumer { |
| 7 |
|
| 8 |
/** |
| 9 |
* @var string the host to connect to (e.g. api.mixpanel.com) |
| 10 |
*/ |
| 11 |
protected $_host; |
| 12 |
|
| 13 |
|
| 14 |
/** |
| 15 |
* @var string the host-relative endpoint to write to (e.g. /engage) |
| 16 |
*/ |
| 17 |
protected $_endpoint; |
| 18 |
|
| 19 |
|
| 20 |
/** |
| 21 |
* @var int connect_timeout The number of seconds to wait while trying to connect. Default is 5 seconds. |
| 22 |
*/ |
| 23 |
protected $_connect_timeout; |
| 24 |
|
| 25 |
|
| 26 |
/** |
| 27 |
* @var int timeout The maximum number of seconds to allow cURL call to execute. Default is 30 seconds. |
| 28 |
*/ |
| 29 |
protected $_timeout; |
| 30 |
|
| 31 |
|
| 32 |
/** |
| 33 |
* @var string the protocol to use for the cURL connection |
| 34 |
*/ |
| 35 |
protected $_protocol; |
| 36 |
|
| 37 |
|
| 38 |
/** |
| 39 |
* @var bool|null true to fork the cURL process (using exec) or false to use PHP's cURL extension. false by default |
| 40 |
*/ |
| 41 |
protected $_fork = null; |
| 42 |
|
| 43 |
|
| 44 |
/** |
| 45 |
* @var int number of cURL requests to run in parallel. 1 by default |
| 46 |
*/ |
| 47 |
protected $_num_threads; |
| 48 |
|
| 49 |
|
| 50 |
/** |
| 51 |
* Creates a new CurlConsumer and assigns properties from the $options array |
| 52 |
* @param array $options |
| 53 |
* @throws Exception |
| 54 |
*/ |
| 55 |
function __construct($options) { |
| 56 |
parent::__construct($options); |
| 57 |
|
| 58 |
$this->_host = $options['host']; |
| 59 |
$this->_endpoint = $options['endpoint']; |
| 60 |
$this->_connect_timeout = isset($options['connect_timeout']) ? $options['connect_timeout'] : 5; |
| 61 |
$this->_timeout = isset($options['timeout']) ? $options['timeout'] : 30; |
| 62 |
$this->_protocol = isset($options['use_ssl']) && $options['use_ssl'] == true ? "https" : "http"; |
| 63 |
$this->_fork = isset($options['fork']) ? ($options['fork'] == true) : false; |
| 64 |
$this->_num_threads = isset($options['num_threads']) ? max(1, intval($options['num_threads'])) : 1; |
| 65 |
|
| 66 |
// ensure the environment is workable for the given settings |
| 67 |
if ($this->_fork == true) { |
| 68 |
$exists = function_exists('exec'); |
| 69 |
if (!$exists) { |
| 70 |
throw new Exception('The "exec" function must exist to use the cURL consumer in "fork" mode. Try setting fork = false or use another consumer.'); |
| 71 |
} |
| 72 |
$disabled = explode(', ', ini_get('disable_functions')); |
| 73 |
$enabled = !in_array('exec', $disabled); |
| 74 |
if (!$enabled) { |
| 75 |
throw new Exception('The "exec" function must be enabled to use the cURL consumer in "fork" mode. Try setting fork = false or use another consumer.'); |
| 76 |
} |
| 77 |
} else { |
| 78 |
if (!function_exists('curl_init')) { |
| 79 |
throw new Exception('The cURL PHP extension is required to use the cURL consumer with fork = false. Try setting fork = true or use another consumer.'); |
| 80 |
} |
| 81 |
} |
| 82 |
} |
| 83 |
|
| 84 |
|
| 85 |
/** |
| 86 |
* Write to the given host/endpoint using either a forked cURL process or using PHP's cURL extension |
| 87 |
* @param array $batch |
| 88 |
* @return bool |
| 89 |
*/ |
| 90 |
public function persist($batch) { |
| 91 |
if (count($batch) > 0) { |
| 92 |
$url = $this->_protocol . "://" . $this->_host . $this->_endpoint; |
| 93 |
if ($this->_fork) { |
| 94 |
$data = "data=" . $this->_encode($batch); |
| 95 |
return $this->_execute_forked($url, $data); |
| 96 |
} else { |
| 97 |
return $this->_execute($url, $batch); |
| 98 |
} |
| 99 |
} else { |
| 100 |
return true; |
| 101 |
} |
| 102 |
} |
| 103 |
|
| 104 |
|
| 105 |
/** |
| 106 |
* Write using the cURL php extension |
| 107 |
* @param $url |
| 108 |
* @param $batch |
| 109 |
* @return bool |
| 110 |
*/ |
| 111 |
protected function _execute($url, $batch) { |
| 112 |
if ($this->_debug()) { |
| 113 |
$this->_log("Making blocking cURL call to $url"); |
| 114 |
} |
| 115 |
|
| 116 |
$mh = curl_multi_init(); |
| 117 |
$chs = array(); |
| 118 |
|
| 119 |
$batch_size = ceil(count($batch) / $this->_num_threads); |
| 120 |
for ($i=0; $i<$this->_num_threads && !empty($batch); $i++) { |
| 121 |
$ch = curl_init(); |
| 122 |
$chs[] = $ch; |
| 123 |
$data = "data=" . $this->_encode(array_splice($batch, 0, $batch_size)); |
| 124 |
curl_setopt($ch, CURLOPT_URL, $url); |
| 125 |
curl_setopt($ch, CURLOPT_HEADER, 0); |
| 126 |
curl_setopt($ch, CURLOPT_CONNECTTIMEOUT, $this->_connect_timeout); |
| 127 |
curl_setopt($ch, CURLOPT_TIMEOUT, $this->_timeout); |
| 128 |
curl_setopt($ch, CURLOPT_POST, 1); |
| 129 |
curl_setopt($ch, CURLOPT_RETURNTRANSFER, 1); |
| 130 |
curl_setopt($ch, CURLOPT_POSTFIELDS, $data); |
| 131 |
curl_multi_add_handle($mh,$ch); |
| 132 |
} |
| 133 |
|
| 134 |
$running = 0; |
| 135 |
|
| 136 |
do { |
| 137 |
curl_multi_exec($mh, $running); |
| 138 |
curl_multi_select($mh); |
| 139 |
} while ($running > 0); |
| 140 |
|
| 141 |
$info = curl_multi_info_read($mh); |
| 142 |
|
| 143 |
$error = false; |
| 144 |
foreach ($chs as $ch) { |
| 145 |
$response = curl_multi_getcontent($ch); |
| 146 |
if (false === $response) { |
| 147 |
$this->_handleError(curl_errno($ch), curl_error($ch)); |
| 148 |
$error = true; |
| 149 |
} |
| 150 |
elseif ("1" != trim($response)) { |
| 151 |
$this->_handleError(0, $response); |
| 152 |
$error = true; |
| 153 |
} |
| 154 |
curl_multi_remove_handle($mh, $ch); |
| 155 |
} |
| 156 |
|
| 157 |
if (CURLE_OK != $info['result']) { |
| 158 |
$this->_handleError($info['result'], "cURL error with code=".$info['result']); |
| 159 |
$error = true; |
| 160 |
} |
| 161 |
|
| 162 |
curl_multi_close($mh); |
| 163 |
return !$error; |
| 164 |
} |
| 165 |
|
| 166 |
|
| 167 |
/** |
| 168 |
* Write using a forked cURL process |
| 169 |
* @param $url |
| 170 |
* @param $data |
| 171 |
* @return bool |
| 172 |
*/ |
| 173 |
protected function _execute_forked($url, $data) { |
| 174 |
|
| 175 |
if ($this->_debug()) { |
| 176 |
$this->_log("Making forked cURL call to $url"); |
| 177 |
} |
| 178 |
|
| 179 |
$exec = 'curl -X POST -H "Content-Type: application/x-www-form-urlencoded" -d ' . $data . ' "' . $url . '"'; |
| 180 |
|
| 181 |
if(!$this->_debug()) { |
| 182 |
$exec .= " >/dev/null 2>&1 &"; |
| 183 |
} |
| 184 |
|
| 185 |
exec($exec, $output, $return_var); |
| 186 |
|
| 187 |
if ($return_var != 0) { |
| 188 |
$this->_handleError($return_var, $output); |
| 189 |
} |
| 190 |
|
| 191 |
return $return_var == 0; |
| 192 |
} |
| 193 |
|
| 194 |
/** |
| 195 |
* @return int |
| 196 |
*/ |
| 197 |
public function getConnectTimeout() |
| 198 |
{ |
| 199 |
return $this->_connect_timeout; |
| 200 |
} |
| 201 |
|
| 202 |
/** |
| 203 |
* @return string |
| 204 |
*/ |
| 205 |
public function getEndpoint() |
| 206 |
{ |
| 207 |
return $this->_endpoint; |
| 208 |
} |
| 209 |
|
| 210 |
/** |
| 211 |
* @return bool|null |
| 212 |
*/ |
| 213 |
public function getFork() |
| 214 |
{ |
| 215 |
return $this->_fork; |
| 216 |
} |
| 217 |
|
| 218 |
/** |
| 219 |
* @return string |
| 220 |
*/ |
| 221 |
public function getHost() |
| 222 |
{ |
| 223 |
return $this->_host; |
| 224 |
} |
| 225 |
|
| 226 |
/** |
| 227 |
* @return array |
| 228 |
*/ |
| 229 |
public function getOptions() |
| 230 |
{ |
| 231 |
return $this->_options; |
| 232 |
} |
| 233 |
|
| 234 |
/** |
| 235 |
* @return string |
| 236 |
*/ |
| 237 |
public function getProtocol() |
| 238 |
{ |
| 239 |
return $this->_protocol; |
| 240 |
} |
| 241 |
|
| 242 |
/** |
| 243 |
* @return int |
| 244 |
*/ |
| 245 |
public function getTimeout() |
| 246 |
{ |
| 247 |
return $this->_timeout; |
| 248 |
} |
| 249 |
|
| 250 |
|
| 251 |
/** |
| 252 |
* Number of requests/batches that will be processed in parallel using curl_multi_exec. |
| 253 |
* @return int |
| 254 |
*/ |
| 255 |
public function getNumThreads() { |
| 256 |
return $this->_num_threads; |
| 257 |
} |
| 258 |
} |
| 259 |
|