Skip to content

Commit e2ccbab

Browse files
authored
Merge branch 6.1 to master (#5969)
1 parent b90f531 commit e2ccbab

8 files changed

Lines changed: 279 additions & 17 deletions

ext-src/php_swoole_http.h

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,11 @@
3434
#define SW_ZLIB_ENCODING_GZIP 0x1f
3535
#define SW_ZLIB_ENCODING_DEFLATE 0x0f
3636
#define SW_ZLIB_ENCODING_ANY 0x2f
37+
#if MAX_MEM_LEVEL >= 8
38+
#define SW_ZLIB_DEF_MEM_LEVEL 8
39+
#else
40+
#define SW_ZLIB_DEF_MEM_LEVEL MAX_MEM_LEVEL
41+
#endif
3742
#endif
3843

3944
#ifdef SW_HAVE_BROTLI
@@ -170,7 +175,7 @@ struct Context {
170175
std::shared_ptr<String> zlib_buffer;
171176
#endif
172177

173-
std::shared_ptr<String> frame_buffer;
178+
std::shared_ptr<String> continue_frame_buffer;
174179
WebSocketSettings websocket_settings;
175180

176181
Request request;

ext-src/swoole_http_client_coro.cc

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -134,7 +134,7 @@ class Client {
134134
* allowing access to the sent Request data even after the connection has been closed.
135135
*/
136136
String *tmp_write_buffer = nullptr;
137-
std::shared_ptr<String> frame_buffer;
137+
std::shared_ptr<String> continue_frame_buffer;
138138
bool connection_close = false;
139139
bool completed = false;
140140
bool event_stream = false;
@@ -1486,7 +1486,7 @@ bool Client::recv_response(double timeout) {
14861486
}
14871487

14881488
void Client::recv_websocket_frame(zval *return_value, double timeout) {
1489-
WebSocket::recv_frame(websocket_settings, frame_buffer, socket, return_value, timeout);
1489+
WebSocket::recv_frame(websocket_settings, continue_frame_buffer, socket, return_value, timeout);
14901490
if (ZVAL_IS_EMPTY_STRING(return_value)) {
14911491
close();
14921492
return;

ext-src/swoole_http_response.cc

Lines changed: 22 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -1300,10 +1300,10 @@ ssize_t WebSocket::send_frame(const swoole::WebSocketSettings &settings,
13001300
* return_value is false means socket.read() method returning -1.
13011301
* return_value is null means other error.
13021302
* return_value is empry string means socket is closed.
1303-
* the opcode is returned so the caller can decide when to release the frame_buffer.
1303+
* the opcode is returned so the caller can decide when to release the continue_frame_buffer.
13041304
*/
13051305
void WebSocket::recv_frame(const WebSocketSettings &settings,
1306-
std::shared_ptr<String> &frame_buffer,
1306+
std::shared_ptr<String> &continue_frame_buffer,
13071307
SocketImpl *sock,
13081308
zval *return_value,
13091309
double timeout) {
@@ -1364,39 +1364,47 @@ void WebSocket::recv_frame(const WebSocketSettings &settings,
13641364
}
13651365

13661366
if (opcode == WebSocket::OPCODE_CONTINUATION) {
1367-
if (sw_unlikely(!frame_buffer)) {
1367+
if (sw_unlikely(!continue_frame_buffer)) {
13681368
swoole_warning("A continuation frame cannot stand alone and MUST be preceded by an initial frame whose "
13691369
"opcode indicates either text or binary data.");
13701370
RETURN_NULL();
13711371
}
13721372

13731373
if (sw_likely(frame.payload)) {
1374-
frame_buffer->append(frame.payload, frame.payload_length);
1374+
continue_frame_buffer->append(frame.payload, frame.payload_length);
13751375
}
13761376

13771377
if (frame.header.FIN) {
13781378
uchar complete_opcode = 0;
13791379
uchar complete_flags = 0;
1380-
WebSocket::parse_ext_flags(frame_buffer->offset, &complete_opcode, &complete_flags);
1380+
WebSocket::parse_ext_flags(continue_frame_buffer->offset, &complete_opcode, &complete_flags);
13811381

13821382
if (complete_flags & WebSocket::FLAG_RSV1) {
1383-
if (sw_unlikely(!FrameObject::uncompress(&zpayload, frame_buffer->str, frame_buffer->length))) {
1383+
if (sw_unlikely(!FrameObject::uncompress(
1384+
&zpayload, continue_frame_buffer->str, continue_frame_buffer->length))) {
1385+
continue_frame_buffer.reset();
13841386
swoole_set_last_error(SW_ERROR_PROTOCOL_ERROR);
13851387
RETURN_NULL();
13861388
}
13871389
} else {
1388-
zend::assign_zend_string_by_val(&zpayload, frame_buffer->str, frame_buffer->length);
1390+
zend::assign_zend_string_by_val(
1391+
&zpayload, continue_frame_buffer->str, continue_frame_buffer->length);
13891392
Z_TRY_ADDREF(zpayload);
1390-
frame_buffer.reset();
13911393
}
13921394

1395+
sw_set_bit(complete_flags, WebSocket::FLAG_FIN);
13931396
WebSocket::construct_frame(return_value, complete_opcode, &zpayload, complete_flags);
13941397
zend::object_set(return_value, ZEND_STRL("fd"), sock->get_fd());
13951398
zval_ptr_dtor(&zpayload);
1399+
/**
1400+
* The final frame of the continuous frame sequence has been received,
1401+
* and the complete message has been assembled. Memory can be released immediately.
1402+
*/
1403+
continue_frame_buffer.reset();
13961404
return;
13971405
}
13981406
} else {
1399-
if (sw_unlikely(frame_buffer)) {
1407+
if (sw_unlikely(continue_frame_buffer)) {
14001408
swoole_warning("All fragments of a message, except for the initial frame, must use the continuation "
14011409
"frame opcode(0).");
14021410
RETURN_NULL();
@@ -1415,12 +1423,12 @@ void WebSocket::recv_frame(const WebSocketSettings &settings,
14151423
zval_ptr_dtor(&zpayload);
14161424
return;
14171425
} else {
1418-
frame_buffer = std::make_shared<String>(
1426+
continue_frame_buffer = std::make_shared<String>(
14191427
(frame.payload_length > 0 ? frame.payload_length : SW_WEBSOCKET_DEFAULT_BUFFER),
14201428
sw_zend_string_allocator());
1421-
frame_buffer->offset = WebSocket::get_ext_flags(frame.header.OPCODE, frame.get_flags());
1429+
continue_frame_buffer->offset = WebSocket::get_ext_flags(frame.header.OPCODE, frame.get_flags());
14221430
if (sw_likely(frame.payload)) {
1423-
frame_buffer->append(frame.payload, frame.payload_length);
1431+
continue_frame_buffer->append(frame.payload, frame.payload_length);
14241432
}
14251433
}
14261434
}
@@ -1445,7 +1453,8 @@ static PHP_METHOD(swoole_http_response, recv) {
14451453
Z_PARAM_DOUBLE(timeout)
14461454
ZEND_PARSE_PARAMETERS_END_EX(RETURN_FALSE);
14471455

1448-
WebSocket::recv_frame(ctx->websocket_settings, ctx->frame_buffer, ctx->get_co_socket(), return_value, timeout);
1456+
WebSocket::recv_frame(
1457+
ctx->websocket_settings, ctx->continue_frame_buffer, ctx->get_co_socket(), return_value, timeout);
14491458
if (ZVAL_IS_EMPTY_STRING(return_value)) {
14501459
ctx->close(ctx);
14511460
return;

ext-src/swoole_websocket_server.cc

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -154,6 +154,7 @@ bool FrameObject::pack(String *buffer) {
154154
sw_set_bit(flags, WebSocket::FLAG_FIN);
155155
need_compress = false;
156156
} else if (opcode == WebSocket::OPCODE_CONTINUATION || !(flags & WebSocket::FLAG_FIN)) {
157+
// Continuous frames and WebSocket message frames without the FLAG_FIN flag do not require compression.
157158
need_compress = false;
158159
}
159160

@@ -407,7 +408,7 @@ bool WebSocket::message_compress(String *buffer, const char *data, size_t length
407408
zstream.zalloc = php_zlib_alloc;
408409
zstream.zfree = php_zlib_free;
409410

410-
status = deflateInit2(&zstream, level, Z_DEFLATED, SW_ZLIB_ENCODING_RAW, MAX_MEM_LEVEL, Z_DEFAULT_STRATEGY);
411+
status = deflateInit2(&zstream, level, Z_DEFLATED, SW_ZLIB_ENCODING_RAW, SW_ZLIB_DEF_MEM_LEVEL, Z_DEFAULT_STRATEGY);
411412
if (status != Z_OK) {
412413
php_swoole_fatal_error(E_WARNING, "deflateInit2() failed, Error: [%d]", status);
413414
return false;
Lines changed: 58 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,58 @@
1+
--TEST--
2+
swoole_http_client_coro/websocket: continue frame finish flag
3+
--SKIPIF--
4+
<?php require __DIR__ . '/../../include/skipif.inc';?>
5+
--FILE--
6+
<?php
7+
require __DIR__ . '/../../include/bootstrap.php';
8+
use Swoole\WebSocket\Server;
9+
use Swoole\Coroutine\Http\Client;
10+
use Swoole\WebSocket\Frame;
11+
use SwooleTest\ProcessManager as ProcessManager;
12+
13+
$data1 = bin2hex(random_bytes(10 * 1024));
14+
$data2 = bin2hex(random_bytes(20 * 2048));
15+
$data3 = bin2hex(random_bytes(40 * 4096));
16+
17+
$pm = new ProcessManager;
18+
$pm->parentFunc = function (int $pid) use ($pm, $data1, $data2, $data3) {
19+
Co\run(function () use ($pm, $data1, $data2, $data3) {
20+
$client = new Client('127.0.0.1', $pm->getFreePort());
21+
$ret = $client->upgrade('/');
22+
Assert::assert($ret);
23+
$client->push($data1, SWOOLE_WEBSOCKET_OPCODE_TEXT, 0);
24+
$client->push($data2, SWOOLE_WEBSOCKET_OPCODE_CONTINUATION, 0);
25+
$client->push($data3, SWOOLE_WEBSOCKET_OPCODE_CONTINUATION, SWOOLE_WEBSOCKET_FLAG_FIN);
26+
$frame = $client->recv();
27+
Assert::true($frame->data == $data1 . $data2 . $data3);
28+
Assert::eq($frame->opcode, SWOOLE_WEBSOCKET_OPCODE_TEXT);
29+
Assert::eq($frame->finish, true);
30+
});
31+
$pm->kill();
32+
};
33+
34+
$pm->childFunc = function () use ($pm, $data1, $data2, $data3) {
35+
$server = new Server('127.0.0.1', $pm->getFreePort(), SERVER_MODE_RANDOM);
36+
$server->set([
37+
'package_max_length' => 100 * 1024 * 1024,
38+
]);
39+
40+
$server->on('workerStart', function () use ($pm) {
41+
$pm->wakeup();
42+
});
43+
44+
$server->on('message', function (Server $server, Frame $frame) use ($pm, $data1, $data2, $data3) {
45+
Assert::true($frame->data == $data1 . $data2 . $data3);
46+
Assert::eq($frame->opcode, SWOOLE_WEBSOCKET_OPCODE_TEXT);
47+
Assert::eq($frame->finish, true);
48+
$server->push($frame->fd, $data1, SWOOLE_WEBSOCKET_OPCODE_TEXT, 0);
49+
$server->push($frame->fd, $data2, SWOOLE_WEBSOCKET_OPCODE_CONTINUATION, 0);
50+
$server->push($frame->fd, $data3, SWOOLE_WEBSOCKET_OPCODE_CONTINUATION, SWOOLE_WEBSOCKET_FLAG_FIN);
51+
});
52+
53+
$server->start();
54+
};
55+
$pm->childFirst();
56+
$pm->run();
57+
?>
58+
--EXPECT--
Lines changed: 54 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,54 @@
1+
--TEST--
2+
swoole_http_client_coro/websocket: continue frame finish flag - 2
3+
--SKIPIF--
4+
<?php require __DIR__ . '/../../include/skipif.inc';?>
5+
--FILE--
6+
<?php
7+
require __DIR__ . '/../../include/bootstrap.php';
8+
use Swoole\Coroutine\Http\Server;
9+
use Swoole\Coroutine\Http\Client;
10+
use Swoole\WebSocket\Frame;
11+
use SwooleTest\ProcessManager as ProcessManager;
12+
13+
$data1 = bin2hex(random_bytes(10 * 1024));
14+
$data2 = bin2hex(random_bytes(20 * 2048));
15+
$data3 = bin2hex(random_bytes(40 * 4096));
16+
17+
$pm = new ProcessManager;
18+
$pm->parentFunc = function (int $pid) use ($pm, $data1, $data2, $data3) {
19+
Co\run(function () use ($pm, $data1, $data2, $data3) {
20+
$client = new Client('127.0.0.1', $pm->getFreePort());
21+
$ret = $client->upgrade('/');
22+
Assert::assert($ret);
23+
$client->push($data1, SWOOLE_WEBSOCKET_OPCODE_TEXT, 0);
24+
$client->push($data2, SWOOLE_WEBSOCKET_OPCODE_CONTINUATION, 0);
25+
$client->push($data3, SWOOLE_WEBSOCKET_OPCODE_CONTINUATION, SWOOLE_WEBSOCKET_FLAG_FIN);
26+
$frame = $client->recv();
27+
Assert::true($frame->data == $data1 . $data2 . $data3);
28+
Assert::eq($frame->opcode, SWOOLE_WEBSOCKET_OPCODE_TEXT);
29+
Assert::eq($frame->finish, true);
30+
});
31+
$pm->kill();
32+
};
33+
34+
$pm->childFunc = function () use ($pm, $data1, $data2, $data3) {
35+
Co\run(function () use ($pm, $data1, $data2, $data3) {
36+
$server = new Server("127.0.0.1", $pm->getFreePort(), false);
37+
$server->handle('/', function ($request, $response) use ($data1, $data2, $data3) {
38+
$response->upgrade();
39+
$frame = $response->recv();
40+
Assert::true($frame->data == $data1 . $data2 . $data3);
41+
Assert::eq($frame->opcode, SWOOLE_WEBSOCKET_OPCODE_TEXT);
42+
Assert::eq($frame->finish, true);
43+
$response->push($data1, SWOOLE_WEBSOCKET_OPCODE_TEXT, 0);
44+
$response->push($data2, SWOOLE_WEBSOCKET_OPCODE_CONTINUATION, 0);
45+
$response->push($data3, SWOOLE_WEBSOCKET_OPCODE_CONTINUATION, SWOOLE_WEBSOCKET_FLAG_FIN);
46+
});
47+
$pm->wakeup();
48+
$server->start();
49+
});
50+
};
51+
$pm->childFirst();
52+
$pm->run();
53+
?>
54+
--EXPECT--
Lines changed: 70 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,70 @@
1+
--TEST--
2+
swoole_http_client_coro/websocket: send more continue frames - websocket server
3+
--SKIPIF--
4+
<?php require __DIR__ . '/../../include/skipif.inc';?>
5+
--FILE--
6+
<?php
7+
require __DIR__ . '/../../include/bootstrap.php';
8+
use Swoole\WebSocket\Server;
9+
use Swoole\Coroutine\Http\Client;
10+
use Swoole\WebSocket\Frame;
11+
use SwooleTest\ProcessManager as ProcessManager;
12+
13+
$data1 = bin2hex(random_bytes(10 * 1024));
14+
$data2 = bin2hex(random_bytes(20 * 2048));
15+
$data3 = bin2hex(random_bytes(40 * 4096));
16+
17+
$pm = new ProcessManager;
18+
$pm->parentFunc = function (int $pid) use ($pm, $data1, $data2, $data3) {
19+
Co\run(function () use ($pm, $data1, $data2, $data3) {
20+
$results = [];
21+
$client = new Client('127.0.0.1', $pm->getFreePort());
22+
$client->set(['websocket_compression' => true]);
23+
$ret = $client->upgrade('/');
24+
Assert::assert($ret);
25+
$client->push('111', SWOOLE_WEBSOCKET_OPCODE_TEXT, SWOOLE_WEBSOCKET_FLAG_FIN);
26+
$frame = $client->recv();
27+
Assert::true($frame->data == $data1 . $data2 . $data2 . $data2 . $data3);
28+
$frame = $client->recv();
29+
Assert::true($frame->data == $data2 . $data1 . $data3);
30+
$frame = $client->recv();
31+
Assert::true($frame->data == $data3 . $data2 . $data1);
32+
});
33+
$pm->kill();
34+
};
35+
36+
$pm->childFunc = function () use ($pm, $data1, $data2, $data3) {
37+
$server = new Server('127.0.0.1', $pm->getFreePort(), SWOOLE_BASE);
38+
$server->set([
39+
'worker_num' => 1,
40+
'package_max_length' => 100 * 1024 * 1024,
41+
'websocket_compression' => true
42+
]);
43+
44+
$server->on('workerStart', function () use ($pm) {
45+
$pm->wakeup();
46+
});
47+
48+
$server->on('message', function (Server $server, Frame $frame) use ($pm, $data1, $data2, $data3) {
49+
$context = deflate_init(ZLIB_ENCODING_RAW);
50+
$server->push($frame->fd, deflate_add($context, $data1, ZLIB_SYNC_FLUSH), SWOOLE_WEBSOCKET_OPCODE_TEXT, SWOOLE_WEBSOCKET_FLAG_COMPRESS | SWOOLE_WEBSOCKET_FLAG_RSV1);
51+
$server->push($frame->fd, deflate_add($context, $data2, ZLIB_SYNC_FLUSH), SWOOLE_WEBSOCKET_OPCODE_CONTINUATION, 0);
52+
$server->push($frame->fd, deflate_add($context, $data2, ZLIB_SYNC_FLUSH), SWOOLE_WEBSOCKET_OPCODE_CONTINUATION, 0);
53+
$server->push($frame->fd, deflate_add($context, $data2, ZLIB_SYNC_FLUSH), SWOOLE_WEBSOCKET_OPCODE_CONTINUATION, 0);
54+
$server->push($frame->fd, deflate_add($context, $data3, ZLIB_FINISH), SWOOLE_WEBSOCKET_OPCODE_CONTINUATION, SWOOLE_WEBSOCKET_FLAG_FIN);
55+
56+
$server->push($frame->fd, deflate_add($context, $data2, ZLIB_NO_FLUSH), SWOOLE_WEBSOCKET_OPCODE_TEXT, SWOOLE_WEBSOCKET_FLAG_COMPRESS | SWOOLE_WEBSOCKET_FLAG_RSV1);
57+
$server->push($frame->fd, deflate_add($context, $data1, ZLIB_NO_FLUSH), SWOOLE_WEBSOCKET_OPCODE_CONTINUATION, 0);
58+
$server->push($frame->fd, deflate_add($context, $data3, ZLIB_FINISH), SWOOLE_WEBSOCKET_OPCODE_CONTINUATION, SWOOLE_WEBSOCKET_FLAG_FIN);
59+
60+
$server->push($frame->fd, deflate_add($context, $data3, ZLIB_NO_FLUSH), SWOOLE_WEBSOCKET_OPCODE_TEXT, SWOOLE_WEBSOCKET_FLAG_COMPRESS | SWOOLE_WEBSOCKET_FLAG_RSV1);
61+
$server->push($frame->fd, deflate_add($context, $data2, ZLIB_NO_FLUSH), SWOOLE_WEBSOCKET_OPCODE_CONTINUATION, 0);
62+
$server->push($frame->fd, deflate_add($context, $data1, ZLIB_FINISH), SWOOLE_WEBSOCKET_OPCODE_CONTINUATION, SWOOLE_WEBSOCKET_FLAG_FIN);
63+
});
64+
65+
$server->start();
66+
};
67+
$pm->childFirst();
68+
$pm->run();
69+
?>
70+
--EXPECT--

0 commit comments

Comments
 (0)