fix(transport): Zstd-распаковка, гонка PacketReceiver, fix имён и стикеров

This commit is contained in:
klockky
2026-05-14 13:12:41 +03:00
parent ea3042f341
commit d541e9df9d
10 changed files with 137 additions and 40 deletions
+40 -15
View File
@@ -1,6 +1,7 @@
import 'dart:typed_data';
import 'dart:isolate';
import 'package:dart_lz4/dart_lz4.dart';
import 'package:es_compression/zstd.dart';
import 'package:msgpack_dart/msgpack_dart.dart' as msgpack;
/// ver(1) + cmd(1) + seq(2) + opcode(2) + packedLen(4) = 10
@@ -129,21 +130,7 @@ Future<Packet> unpackPacket(Uint8List packet) async {
if (payloadBytes.isNotEmpty) {
if (compFlag != 0) {
try {
payloadBytes = lz4Decompress(
payloadBytes,
decompressedSize: _maxDecompressedSize,
);
} catch (_) {
try {
payloadBytes = _lz4BlockDecompress(
payloadBytes,
_maxDecompressedSize,
);
} catch (e) {
throw Exception("LZ4 decompression error: $e");
}
}
payloadBytes = _decompressPayload(payloadBytes);
}
try {
@@ -165,6 +152,44 @@ Future<Packet> unpackPacket(Uint8List packet) async {
});
}
/// Определяет формат сжатия по magic-number и распаковывает payload.
/// Сервер может присылать LZ4 block ИЛИ Zstandard в зависимости от ответа.
Uint8List _decompressPayload(Uint8List src) {
// Zstandard: magic 28 B5 2F FD (little-endian)
if (src.length >= 4 &&
src[0] == 0x28 &&
src[1] == 0xB5 &&
src[2] == 0x2F &&
src[3] == 0xFD) {
try {
final out = zstd.decode(src);
return out is Uint8List ? out : Uint8List.fromList(out);
} catch (e) {
throw Exception('Zstd decompression error: $e');
}
}
// LZ4 frame: magic 04 22 4D 18
if (src.length >= 4 &&
src[0] == 0x04 &&
src[1] == 0x22 &&
src[2] == 0x4D &&
src[3] == 0x18) {
try {
return lz4Decompress(src, decompressedSize: _maxDecompressedSize);
} catch (e) {
throw Exception('LZ4 frame decompression error: $e');
}
}
// По умолчанию — LZ4 block (без magic)
try {
return _lz4BlockDecompress(src, _maxDecompressedSize);
} catch (e) {
throw Exception('LZ4 block decompression error: $e');
}
}
/// LZ4 block декомпрессия (без frame-заголовка).
/// Сервер шлёт именно block-формат, dart_lz4 его не поддерживает.
Uint8List _lz4BlockDecompress(Uint8List src, int maxSize) {
+9 -12
View File
@@ -4,15 +4,16 @@ import '../protocol/packet.dart';
import '../utils/logger.dart';
/// Буфер входящих данных.
/// Копит сырые байты из сокета, собирает из них целые пакеты.
/// Копит сырые байты из сокета, нарезает их на байтовые срезы целых пакетов.
class PacketReceiver {
Uint8List _buffer = Uint8List(0);
static const int _maxBufferSize = 2 * 1024 * 1024; // 2 мегабуйта
/// Добавляет байты в буфер, возвращает поток собранных пакетов.
/// Неполные данные остаются в буфере до следующего вызова.
Stream<Packet> feed(Uint8List data) async* {
/// Добавляет байты в буфер и возвращает все собранные пакеты как сырые срезы.
/// Полностью синхронный — нарезка не блокируется на распаковке, поэтому
/// конкурентные вызовы из stream-листенера не могут пересечься на `_buffer`.
List<Uint8List> feed(Uint8List data) {
final newBuffer = Uint8List(_buffer.length + data.length);
newBuffer.setAll(0, _buffer);
newBuffer.setAll(_buffer.length, data);
@@ -23,9 +24,10 @@ class PacketReceiver {
'PacketReceiver: переполнение буфера (${_buffer.length} B), сброс',
);
reset();
return;
return const [];
}
final packets = <Uint8List>[];
while (_buffer.length >= headerSize) {
final bd = ByteData.view(
_buffer.buffer,
@@ -38,15 +40,10 @@ class PacketReceiver {
if (_buffer.length < totalLength) break;
final packetBytes = Uint8List.sublistView(_buffer, 0, totalLength);
packets.add(Uint8List.sublistView(_buffer, 0, totalLength));
_buffer = _buffer.sublist(totalLength);
try {
yield await unpackPacket(packetBytes);
} catch (e) {
logger.e('PacketReceiver: ошибка распаковки: $e');
}
}
return packets;
}
void reset() {