PluginProbe
Imagify Image Optimization: Optimize Images | Compress & Convert to WebP/AVIF / trunk
Imagify Image Optimization: Optimize Images | Compress & Convert to WebP/AVIF vtrunk
2.3.4 2.3.3 2.3.2 2.3.1 2.3.0 2.2.9 2.2.8 trunk 1.10 1.3.3 1.3.4 1.3.5 1.3.5.1 1.3.5.2 1.3.6 1.3.6.1 1.4 1.4.1 1.4.2 1.4.3 1.4.4 1.4.5 1.4.6 1.4.7 1.5 All 103 releases
imagify / vendor / wp-media / wp-mixpanel / src / Classes / ConsumerStrategies / SocketConsumer.php

SocketConsumer.php in Imagify Image Optimization: Optimize Images | Compress & Convert to WebP/AVIF trunk, at vendor/wp-media/wp-mixpanel/src/Classes/ConsumerStrategies/SocketConsumer.php

314 lines 9.0 KB
No matching file
Up and down to move Enter to open Esc to close
Raw Download Zip
1 <?php
2 /**
3 * Portions of this class were borrowed from
4 * https://github.com/segmentio/analytics-php/blob/master/lib/Analytics/Consumer/Socket.php.
5 * Thanks for the work!
6 *
7 * WWWWWW||WWWWWW
8 * W W W||W W W
9 * ||
10 * ( OO )__________
11 * / | \
12 * /o o| MIT \
13 * \___/||_||__||_|| *
14 * || || || ||
15 * _||_|| _||_||
16 * (__|__|(__|__|
17 * (The MIT License)
18 *
19 * Copyright (c) 2013 Segment.io Inc. [email protected]
20 *
21 * Permission is hereby granted, free of charge, to any person obtaining a copy of this software and associated
22 * documentation files (the 'Software'), to deal in the Software without restriction, including without limitation the
23 * rights to use, copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the Software, and to
24 * permit persons to whom the Software is furnished to do so, subject to the following conditions:
25 *
26 * The above copyright notice and this permission notice shall be included in all copies or substantial portions of the
27 * Software.
28 *
29 * THE SOFTWARE IS PROVIDED 'AS IS', WITHOUT WARRANTY OF ANY KIND, EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE
30 * WARRANTIES OF MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS
31 * OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR
32 * OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
33 */
34 require_once(dirname(__FILE__) . "/AbstractConsumer.php");
35
36 /**
37 * Consumes messages and writes them to host/endpoint using a persistent socket
38 */
39 class Imagify_WPMedia_ConsumerStrategies_SocketConsumer extends Imagify_WPMedia_ConsumerStrategies_AbstractConsumer {
40
41 /**
42 * @var string the host to connect to (e.g. api.mixpanel.com)
43 */
44 private $_host;
45
46
47 /**
48 * @var string the host-relative endpoint to write to (e.g. /engage)
49 */
50 private $_endpoint;
51
52
53 /**
54 * @var int connect_timeout the socket connection timeout in seconds
55 */
56 private $_connect_timeout;
57
58
59 /**
60 * @var string the protocol to use for the socket connection
61 */
62 private $_protocol;
63
64
65 /**
66 * @var resource holds the socket resource
67 */
68 private $_socket;
69
70 /**
71 * @var bool whether or not to wait for a response
72 */
73 private $_async;
74
75 /**
76 * @var int the port to use for the socket connection
77 */
78 private $_port;
79
80
81 /**
82 * Creates a new SocketConsumer and assigns properties from the $options array
83 * @param array $options
84 */
85 public function __construct($options = array()) {
86 parent::__construct($options);
87
88
89 $this->_host = $options['host'];
90 $this->_endpoint = $options['endpoint'];
91 $this->_connect_timeout = isset($options['connect_timeout']) ? $options['connect_timeout'] : 5;
92 $this->_async = isset($options['async']) && $options['async'] === false ? false : true;
93
94 if (array_key_exists('use_ssl', $options) && $options['use_ssl'] == true) {
95 $this->_protocol = "ssl";
96 $this->_port = 443;
97 } else {
98 $this->_protocol = "tcp";
99 $this->_port = 80;
100 }
101 }
102
103
104 /**
105 * Write using a persistent socket connection.
106 * @param array $batch
107 * @return bool
108 */
109 public function persist($batch) {
110
111 $socket = $this->_getSocket();
112 if (!is_resource($socket)) {
113 return false;
114 }
115
116 $data = "data=".$this->_encode($batch);
117
118 $body = "";
119 $body.= "POST ".$this->_endpoint." HTTP/1.1\r\n";
120 $body.= "Host: " . $this->_host . "\r\n";
121 $body.= "Content-Type: application/x-www-form-urlencoded\r\n";
122 $body.= "Accept: application/json\r\n";
123 $body.= "Content-length: " . strlen($data) . "\r\n";
124 $body.= "\r\n";
125 $body.= $data;
126
127 return $this->_write($socket, $body);
128 }
129
130
131 /**
132 * Return cached socket if open or create a new persistent socket
133 * @return bool|resource
134 */
135 private function _getSocket() {
136 if(is_resource($this->_socket)) {
137
138 if ($this->_debug()) {
139 $this->_log("Using existing socket");
140 }
141
142 return $this->_socket;
143 } else {
144
145 if ($this->_debug()) {
146 $this->_log("Creating new socket at ".time());
147 }
148
149 return $this->_createSocket();
150 }
151 }
152
153 /**
154 * Attempt to open a new socket connection, cache it, and return the resource
155 * @param bool $retry
156 * @return bool|resource
157 */
158 private function _createSocket($retry = true) {
159 try {
160 $socket = pfsockopen($this->_protocol . "://" . $this->_host, $this->_port, $err_no, $err_msg, $this->_connect_timeout);
161
162 if ($this->_debug()) {
163 $this->_log("Opening socket connection to " . $this->_protocol . "://" . $this->_host . ":" . $this->_port);
164 }
165
166 if ($err_no != 0) {
167 $this->_handleError($err_no, $err_msg);
168 return $retry == true ? $this->_createSocket(false) : false;
169 } else {
170 // cache the socket
171 $this->_socket = $socket;
172 return $socket;
173 }
174
175 } catch (Exception $e) {
176 $this->_handleError($e->getCode(), $e->getMessage());
177 return $retry == true ? $this->_createSocket(false) : false;
178 }
179 }
180
181 /**
182 * Attempt to close and dereference a socket resource
183 */
184 private function _destroySocket() {
185 $socket = $this->_socket;
186 $this->_socket = null;
187 fclose($socket);
188 }
189
190
191 /**
192 * Write $data through the given $socket
193 * @param $socket
194 * @param $data
195 * @param bool $retry
196 * @return bool
197 */
198 private function _write($socket, $data, $retry = true) {
199
200 $bytes_sent = 0;
201 $bytes_total = strlen($data);
202 $socket_closed = false;
203 $success = true;
204 $max_bytes_per_write = 8192;
205
206 // if we have no data to write just return true
207 if ($bytes_total == 0) {
208 return true;
209 }
210
211 // try to write the data
212 while (!$socket_closed && $bytes_sent < $bytes_total) {
213
214 try {
215 $bytes = fwrite($socket, $data, $max_bytes_per_write);
216
217 if ($this->_debug()) {
218 $this->_log("Socket wrote ".$bytes." bytes");
219 }
220
221 // if we actually wrote data, then remove the written portion from $data left to write
222 if ($bytes > 0) {
223 $data = substr($data, $max_bytes_per_write);
224 }
225
226 } catch (Exception $e) {
227 $this->_handleError($e->getCode(), $e->getMessage());
228 $socket_closed = true;
229 }
230
231 if (isset($bytes) && $bytes) {
232 $bytes_sent += $bytes;
233 } else {
234 $socket_closed = true;
235 }
236 }
237
238 // create a new socket if the current one is closed and retry the message
239 if ($socket_closed) {
240
241 $this->_destroySocket();
242
243 if ($retry) {
244 if ($this->_debug()) {
245 $this->_log("Retrying socket write...");
246 }
247 $socket = $this->_getSocket();
248 if ($socket) return $this->_write($socket, $data, false);
249 }
250
251 return false;
252 }
253
254
255 // only wait for the response in debug mode or if we explicitly want to be synchronous
256 if ($this->_debug() || !$this->_async) {
257 $res = $this->handleResponse(fread($socket, 2048));
258 if ($res["status"] != "200") {
259 $this->_handleError($res["status"], $res["body"]);
260 $success = false;
261 }
262 }
263
264 return $success;
265 }
266
267
268 /**
269 * Parse the response from a socket write (only used for debugging)
270 * @param $response
271 * @return array
272 */
273 private function handleResponse($response) {
274
275 $lines = explode("\n", $response);
276
277 // extract headers
278 $headers = array();
279 foreach($lines as $line) {
280 $kvsplit = explode(":", $line);
281 if (count($kvsplit) == 2) {
282 $header = $kvsplit[0];
283 $value = $kvsplit[1];
284 $headers[$header] = trim($value);
285 }
286
287 }
288
289 // extract status
290 $line_one_exploded = explode(" ", $lines[0]);
291 $status = $line_one_exploded[1];
292
293 // extract body
294 $body = $lines[count($lines) - 1];
295
296 // if the connection has been closed lets kill the socket
297 if (isset($headers["Connection"]) and $headers['Connection'] == "close") {
298 $this->_destroySocket();
299 if ($this->_debug()) {
300 $this->_log("Server told us connection closed so lets destroy the socket so it'll reconnect on next call");
301 }
302 }
303
304 $ret = array(
305 "status" => $status,
306 "body" => $body,
307 );
308
309 return $ret;
310 }
311
312
313 }
314