@@ -297,6 +297,7 @@ bool next_control_message(std::span<const std::uint8_t> bytes, std::size_t& mess
297297 }
298298
299299 switch (type) {
300+ // uint16 payload length: type varint + uint16 length + payload
300301 case kClientSetupType :
301302 case kServerSetupType :
302303 case kSubscribeUpdateType :
@@ -307,7 +308,19 @@ bool next_control_message(std::span<const std::uint8_t> bytes, std::size_t& mess
307308 case kPublishNamespaceOkType :
308309 case kPublishNamespaceErrorType :
309310 case kPublishNamespaceDoneType :
310- case kPublishDoneType : {
311+ case kPublishDoneType :
312+ case kMaxRequestIdType :
313+ case kPublishType :
314+ case kPublishOkType :
315+ case kPublishErrorType :
316+ case 0x0a : // UNSUBSCRIBE
317+ case 0x10 : // GOAWAY
318+ case 0x14 : // SUBSCRIBE_DONE (draft-14)
319+ case 0x16 : // FETCH
320+ case 0x17 : // FETCH_OK
321+ case 0x18 : // FETCH_CANCEL
322+ case 0x1a : // REQUESTS_BLOCKED (draft-16)
323+ {
311324 if (offset + 2 > bytes.size ()) {
312325 return false ;
313326 }
@@ -316,41 +329,22 @@ bool next_control_message(std::span<const std::uint8_t> bytes, std::size_t& mess
316329 message_size = offset + 2 + payload_length;
317330 return bytes.size () >= message_size;
318331 }
319- case kSubscribeNamespaceType :
320- case kSubscribeNamespaceOkType :
321- case kPublishType :
322- case kPublishOkType :
323- case kPublishErrorType : {
324- if (type == kPublishType || type == kPublishOkType || type == kPublishErrorType ) {
325- if (offset + 2 > bytes.size ()) {
326- return false ;
327- }
328- const std::size_t payload_length =
329- (static_cast <std::size_t >(bytes[offset]) << 8 ) | static_cast <std::size_t >(bytes[offset + 1 ]);
330- message_size = offset + 2 + payload_length;
331- return bytes.size () >= message_size;
332- }
332+ // varint payload length: type varint + varint length + payload
333+ // SUBSCRIBE_NAMESPACE family — same framing as 0x11 and 0x12.
334+ case kSubscribeNamespaceType : // 0x11
335+ case kSubscribeNamespaceOkType : // 0x12
336+ case 0x13 : // SUBSCRIBE_NAMESPACE_ERROR (draft-14)
337+ case 0x1b : // UNSUBSCRIBE_NAMESPACE (draft-14)
338+ {
333339 std::uint64_t payload_length = 0 ;
334340 if (!decode_varint_impl (bytes, offset, payload_length)) {
335341 return false ;
336342 }
337343 message_size = offset + static_cast <std::size_t >(payload_length);
338344 return bytes.size () >= message_size;
339345 }
340- case kMaxRequestIdType :
341- default : {
342- // All known length-prefixed messages (including unknown future types) use
343- // a uint16 payload length immediately after the type varint. Attempt to
344- // consume the message this way so that unrecognised messages (e.g. FETCH,
345- // GOAWAY, UNSUBSCRIBE) do not block the buffer.
346- if (offset + 2 > bytes.size ()) {
347- return false ;
348- }
349- const std::size_t payload_length =
350- (static_cast <std::size_t >(bytes[offset]) << 8 ) | static_cast <std::size_t >(bytes[offset + 1 ]);
351- message_size = offset + 2 + payload_length;
352- return bytes.size () >= message_size;
353- }
346+ default :
347+ return false ;
354348 }
355349}
356350
0 commit comments