Fix/ts splitter resync (#4783)

This commit is contained in:
YuLi
2026-07-22 22:58:31 -07:00
committed by GitHub
parent 4a24057e73
commit 59ffb2c9fd
4 changed files with 59 additions and 12 deletions

View File

@@ -149,6 +149,7 @@ void HttpRequestSplitter::reset() {
_content_len = 0;
_remain_data_size = 0;
_remain_data.clear();
onReset();
}
const char *HttpRequestSplitter::onSearchPacketTail(const char *data,size_t len) {
@@ -179,4 +180,3 @@ HttpRequestSplitter::HttpRequestSplitter() {
}
} /* namespace mediakit */

View File

@@ -68,6 +68,8 @@ public:
void setMaxCacheSize(size_t max_cache_size);
protected:
virtual void onReset() {};
/**
* 收到请求头
* @param data 请求头数据

View File

@@ -19,6 +19,10 @@ void TSSegment::setOnSegment(TSSegment::onSegment cb) {
_onSegment = std::move(cb);
}
void TSSegment::onReset() {
_is_synced = false;
}
ssize_t TSSegment::onRecvHeader(const char *data, size_t len) {
if (!isTSPacket(data, len)) {
WarnL << "不是ts包:" << (int) (data[0]) << " " << len;
@@ -29,22 +33,60 @@ ssize_t TSSegment::onRecvHeader(const char *data, size_t len) {
}
const char *TSSegment::onSearchPacketTail(const char *data, size_t len) {
if (len < _size + 1) {
if (len == _size && ((uint8_t *) data)[0] == TS_SYNC_BYTE) {
return data + _size;
}
auto bytes = (const uint8_t *) data;
if (_is_synced) {
if (len < _size) {
return nullptr;
}
// 下一个包头 [AUTO-TRANSLATED:c653c49d]
// Next packet header
if (((uint8_t *) data)[_size] == TS_SYNC_BYTE) {
// 锁定后只按当前同步字节分包TS 语义和错误处理交给 libmpegts
// Once locked, split only by the current sync byte; leave TS semantics and error handling to libmpegts
if (bytes[0] == TS_SYNC_BYTE) {
return data + _size;
}
auto pos = memchr(data + _size, TS_SYNC_BYTE, len - _size);
if (pos) {
return (char *) pos;
_is_synced = false;
}
// 单个完整包可以直接转交,但不足以确认后续数据的包边界
// A single complete packet can be forwarded, but is insufficient to confirm subsequent packet boundaries
if (len == _size && bytes[0] == TS_SYNC_BYTE) {
return data + _size;
}
return searchPacketTailUnSynced(data, len);
}
#if defined(_MSC_VER)
__declspec(noinline)
#elif defined(__GNUC__)
__attribute__((noinline))
#endif
const char *TSSegment::searchPacketTailUnSynced(const char *data, size_t len) {
if (len > _size) {
const char *candidate = data;
const char *search_end = data + len - _size;
while (candidate < search_end) {
candidate = (const char *) memchr(candidate, TS_SYNC_BYTE, search_end - candidate);
if (!candidate) {
break;
}
if ((uint8_t) candidate[_size] == TS_SYNC_BYTE) {
_is_synced = true;
return candidate == data ? data + _size : candidate;
}
++candidate;
}
}
// 等价于 len > 4 * _size但不会发生无符号整数溢出
// Equivalent to len > 4 * _size without unsigned integer overflow
if (len && _size <= (len - 1) / 4) {
// 尾部同步字节可能在下一次输入后组成同步对,因此只丢弃它之前的数据
// A trailing sync byte may form a sync pair after the next input, so discard only the data before it
const char *tail = data + len - _size;
auto candidate = memchr(tail, TS_SYNC_BYTE, data + len - tail);
if (candidate) {
return (const char *) candidate;
}
if (remainDataSize() > 4 * _size) {
// 数据这么多都没ts包全部清空 [AUTO-TRANSLATED:95bece98]
// So much data but no ts packets, clear all
return data + len;
@@ -90,7 +132,7 @@ TSDecoder::~TSDecoder() {
}
ssize_t TSDecoder::input(const uint8_t *data, size_t bytes) {
if (TSSegment::isTSPacket((char *)data, bytes)) {
if (_ts_segment.remainDataSize() == 0 && TSSegment::isTSPacket((char *)data, bytes)) {
return ts_demuxer_input(_demuxer_ctx, (uint8_t *) data, bytes);
}
try {

View File

@@ -30,11 +30,14 @@ public:
static bool isTSPacket(const char *data, size_t len);
protected:
void onReset() override;
ssize_t onRecvHeader(const char *data, size_t len) override ;
const char *onSearchPacketTail(const char *data, size_t len) override ;
private:
const char *searchPacketTailUnSynced(const char *data, size_t len);
size_t _size;
bool _is_synced = false;
onSegment _onSegment;
};