Skip to content

Commit c9b1248

Browse files
authored
fix(Core): handle oversized messages via E2BIG fallback drain in BatchJob (#9689)
1 parent 07ccec8 commit c9b1248

2 files changed

Lines changed: 102 additions & 1 deletion

File tree

‎Core/src/Batch/BatchJob.php‎

Lines changed: 48 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@ class BatchJob implements JobInterface
3232
const DEFAULT_BATCH_SIZE = 100;
3333
const DEFAULT_CALL_PERIOD = 2.0;
3434
const DEFAULT_WORKERS = 1;
35+
const MAX_MESSAGE_SIZE = 8192;
3536

3637
use JobTrait;
3738
use SysvTrait;
@@ -108,7 +109,7 @@ public function run()
108109
$lastInvoked = microtime(true);
109110
$maxSize = is_array($stat = @msg_stat_queue($q)) && isset($stat['msg_qbytes'])
110111
? $stat['msg_qbytes']
111-
: 8192;
112+
: self::MAX_MESSAGE_SIZE;
112113

113114
if (!is_null($this->bootstrapFile)) {
114115
require_once($this->bootstrapFile);
@@ -133,6 +134,8 @@ public function run()
133134
$items[] = unserialize(file_get_contents($message));
134135
@unlink($message);
135136
}
137+
} elseif ($this->isMsgTooBig($errorcode)) {
138+
$this->drainOversizedMessage($q, $maxSize);
136139
}
137140
pcntl_signal_dispatch();
138141
// It runs the job when
@@ -207,4 +210,48 @@ public function getBatchSize()
207210
{
208211
return $this->batchSize;
209212
}
213+
214+
/**
215+
* Drain an oversized message from the queue to prevent head-of-line blocking.
216+
*
217+
* @access private
218+
* @internal
219+
*
220+
* @param resource $q The message queue resource.
221+
* @param int $maxSize The max buffer size to receive.
222+
* @return bool
223+
*/
224+
public function drainOversizedMessage($q, $maxSize = self::MAX_MESSAGE_SIZE)
225+
{
226+
$discardType = 0;
227+
$discardMessage = null;
228+
$discardErrno = 0;
229+
return @msg_receive(
230+
$q,
231+
0,
232+
$discardType,
233+
$maxSize,
234+
$discardMessage,
235+
false,
236+
MSG_IPC_NOWAIT | MSG_NOERROR,
237+
$discardErrno
238+
);
239+
}
240+
241+
/**
242+
* Check if the error code from msg_receive indicates that the message was too big.
243+
*
244+
* @access private
245+
* @internal
246+
*
247+
* @param int $errorcode
248+
* @return bool
249+
*/
250+
public function isMsgTooBig($errorcode)
251+
{
252+
return (defined('MSG_E2BIG') && $errorcode === constant('MSG_E2BIG'))
253+
|| (defined('PCNTL_E2BIG') && $errorcode === PCNTL_E2BIG)
254+
|| (defined('SOCKET_E2BIG') && $errorcode === SOCKET_E2BIG)
255+
|| $errorcode === 7;
256+
}
210257
}

‎Core/tests/Unit/Batch/BatchJobTest.php‎

Lines changed: 54 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -70,4 +70,58 @@ public function runJob($items)
7070
}
7171
return true;
7272
}
73+
74+
public function testIsMsgTooBig()
75+
{
76+
$job = new BatchJob('testing', array($this, 'runJob'), 1);
77+
$this->assertTrue($job->isMsgTooBig(7));
78+
if (defined('MSG_E2BIG')) {
79+
$this->assertTrue($job->isMsgTooBig(constant('MSG_E2BIG')));
80+
}
81+
if (defined('PCNTL_E2BIG')) {
82+
$this->assertTrue($job->isMsgTooBig(PCNTL_E2BIG));
83+
}
84+
if (defined('SOCKET_E2BIG')) {
85+
$this->assertTrue($job->isMsgTooBig(SOCKET_E2BIG));
86+
}
87+
$this->assertFalse($job->isMsgTooBig(0));
88+
$this->assertFalse($job->isMsgTooBig(4));
89+
$this->assertFalse($job->isMsgTooBig(35));
90+
}
91+
92+
public function testDrainOversizedMessage()
93+
{
94+
if (!extension_loaded('sysvmsg')) {
95+
$this->markTestSkipped('sysvmsg extension required');
96+
}
97+
$key = ftok(__FILE__, 'B');
98+
$q = msg_get_queue($key);
99+
while (@msg_receive($q, 0, $t, 8192, $m, false, MSG_IPC_NOWAIT | MSG_NOERROR, $e)) {
100+
}
101+
$job = new BatchJob('testing', array($this, 'runJob'), 1);
102+
103+
// Verify drain on an empty queue returns false.
104+
$this->assertFalse($job->drainOversizedMessage($q));
105+
106+
// Queue an oversized item and a subsequent valid item.
107+
$oversizedMessage = str_repeat('A', 100);
108+
$nextMessage = 'valid message';
109+
$this->assertTrue(msg_send($q, 1, $oversizedMessage, true, false));
110+
$this->assertTrue(msg_send($q, 1, $nextMessage, true, false));
111+
112+
// Attempting to read with a buffer smaller than the message triggers E2BIG.
113+
$received = @msg_receive($q, 0, $type, 10, $message, true, MSG_IPC_NOWAIT, $errorcode);
114+
$this->assertFalse($received);
115+
$this->assertTrue($job->isMsgTooBig($errorcode));
116+
117+
// Draining with MSG_NOERROR purges the oversized message and unblocks the queue.
118+
$this->assertTrue($job->drainOversizedMessage($q));
119+
120+
// The subsequent message is now at the head of the queue and can be received.
121+
$receivedNext = @msg_receive($q, 0, $type, 8192, $message, true, MSG_IPC_NOWAIT, $errorcode);
122+
$this->assertTrue($receivedNext);
123+
$this->assertSame($nextMessage, $message);
124+
125+
msg_remove_queue($q);
126+
}
73127
}

0 commit comments

Comments
 (0)