7#include <packager/media/chunking/text_chunker.h>
14#include <absl/log/check.h>
15#include <absl/log/log.h>
17#include <packager/chunking_params.h>
18#include <packager/macros/status.h>
19#include <packager/media/base/media_handler.h>
20#include <packager/media/base/stream_info.h>
21#include <packager/media/base/text_sample.h>
22#include <packager/media/base/timestamp_util.h>
23#include <packager/status.h>
28const size_t kStreamIndex = 0;
31TextChunker::TextChunker(
double segment_duration_in_seconds,
32 int64_t start_segment_number)
33 : segment_duration_in_seconds_(segment_duration_in_seconds),
34 segment_number_(start_segment_number),
35 ts_ttx_heartbeat_shift_(kDefaultTtxHeartbeatShift),
36 use_segment_coordinator_(false) {};
38TextChunker::TextChunker(
double segment_duration_in_seconds,
39 int64_t start_segment_number,
40 int64_t ts_ttx_heartbeat_shift)
41 : segment_duration_in_seconds_(segment_duration_in_seconds),
42 segment_number_(start_segment_number),
43 ts_ttx_heartbeat_shift_(ts_ttx_heartbeat_shift),
44 use_segment_coordinator_(false) {};
46TextChunker::TextChunker(
double segment_duration_in_seconds,
47 int64_t start_segment_number,
48 int64_t ts_ttx_heartbeat_shift,
49 bool use_segment_coordinator)
50 : segment_duration_in_seconds_(segment_duration_in_seconds),
51 segment_number_(start_segment_number),
52 ts_ttx_heartbeat_shift_(ts_ttx_heartbeat_shift),
53 use_segment_coordinator_(use_segment_coordinator) {};
55Status TextChunker::Process(std::unique_ptr<StreamData> data) {
56 switch (data->stream_data_type) {
57 case StreamDataType::kStreamInfo:
58 return OnStreamInfo(std::move(data->stream_info));
59 case StreamDataType::kTextSample:
60 return OnTextSample(data->text_sample);
61 case StreamDataType::kCueEvent:
62 return OnCueEvent(data->cue_event);
63 case StreamDataType::kSegmentInfo:
64 if (use_segment_coordinator_) {
65 return OnSegmentInfo(std::move(data->segment_info));
68 return DispatchSegmentInfo(kStreamIndex, std::move(data->segment_info));
71 return Status(error::INTERNAL_ERROR,
72 "Invalid stream data type for this handler");
76Status TextChunker::OnFlushRequest(
size_t ) {
84 while (samples_in_current_segment_.size()) {
85 if (segment_start_ < 0) {
89 int64_t segment_end = segment_start_ + segment_duration_;
90 AddOngoingCuesToCurrentSegment(segment_end);
91 RETURN_IF_ERROR(DispatchSegment(segment_duration_));
94 return FlushAllDownstreams();
97Status TextChunker::OnStreamInfo(std::shared_ptr<const StreamInfo> info) {
98 time_scale_ = info->time_scale();
99 segment_duration_ = ScaleTime(segment_duration_in_seconds_);
101 return DispatchStreamInfo(kStreamIndex, std::move(info));
104Status TextChunker::OnCueEvent(std::shared_ptr<const CueEvent> event) {
113 const int64_t event_time = ScaleTime(event->time_in_seconds);
115 while (segment_start_ + segment_duration_ < event_time) {
116 RETURN_IF_ERROR(DispatchSegment(segment_duration_));
119 const int64_t shorten_duration = event_time - segment_start_;
120 RETURN_IF_ERROR(DispatchSegment(shorten_duration));
121 return DispatchCueEvent(kStreamIndex, std::move(event));
124Status TextChunker::OnTextSample(std::shared_ptr<const TextSample> sample) {
129 int64_t sample_start = sample->start_time();
130 const auto role = sample->role();
131 DVLOG(2) <<
"OnTextSample: role=" <<
static_cast<int>(role)
132 <<
" pts=" << sample_start <<
" end=" << sample->EndTime()
133 <<
" is_empty=" << sample->is_empty()
134 <<
" sub_stream_index=" << sample->sub_stream_index();
139 if (segment_start_ < 0 && !use_segment_coordinator_) {
143 segment_start_ = (sample_start / segment_duration_) * segment_duration_;
144 DVLOG(1) <<
"first segment start=" << segment_start_;
148 case TextSampleRole::kCue: {
149 DVLOG(2) <<
"PTS=" << sample_start <<
" cue with end "
150 << sample->EndTime();
153 case TextSampleRole::kCueStart: {
154 DVLOG(2) <<
"PTS=" << sample_start <<
" cue start wo end";
157 case TextSampleRole::kCueEnd: {
158 DVLOG(2) <<
"PTS=" << sample_start <<
" cue end";
160 auto end_time = sample->EndTime();
161 for (
auto s : samples_without_end_) {
162 int64_t cue_start = s->start_time();
163 if (cue_start < segment_start_) {
164 cue_start = segment_start_;
167 std::make_shared<TextSample>(
"", cue_start, end_time, s->settings(),
168 s->body(), TextSampleRole::kCue);
169 DVLOG(3) <<
"cue shortened. startTime=" << s->start_time()
170 <<
" endTime=" << end_time;
171 samples_in_current_segment_.push_back(nS);
173 samples_without_end_.clear();
176 case TextSampleRole::kTextHeartBeat: {
179 case TextSampleRole::kMediaHeartBeat: {
180 sample_start -= ts_ttx_heartbeat_shift_;
181 latest_media_heartbeat_time_ = sample_start;
182 DVLOG(3) <<
"PTS=" << sample_start <<
" media heartbeat";
190 if (role != TextSampleRole::kMediaHeartBeat) {
192 if (PtsIsBefore(sample_start, latest_media_heartbeat_time_)) {
193 LOG(WARNING) <<
"Potentially bad text segment: text pts=" << sample_start
194 <<
" before latest media pts="
195 << latest_media_heartbeat_time_;
218 if (!use_segment_coordinator_ && segment_start_ >= 0) {
219 int64_t segment_end = segment_start_ + segment_duration_;
220 while (!PtsIsBefore(sample_start, segment_end)) {
222 AddOngoingCuesToCurrentSegment(segment_end);
224 RETURN_IF_ERROR(DispatchSegment(segment_duration_));
225 segment_end = segment_start_ + segment_duration_;
230 case TextSampleRole::kCue: {
231 samples_in_current_segment_.push_back(std::move(sample));
234 case TextSampleRole::kCueStart: {
235 samples_without_end_.push_back(std::move(sample));
246Status TextChunker::OnSegmentInfo(std::shared_ptr<const SegmentInfo> info) {
247 DCHECK(use_segment_coordinator_)
248 <<
"OnSegmentInfo should only be called when coordinator mode is enabled";
251 if (info->is_subsegment) {
252 DVLOG(3) <<
"TextChunker: Skipping subsegment SegmentInfo";
260 int64_t segment_end_boundary = info->start_timestamp + info->duration;
262 DVLOG(2) <<
"TextChunker received SegmentInfo: start="
263 << info->start_timestamp <<
" duration=" << info->duration
264 <<
" end_boundary=" << segment_end_boundary
265 <<
" (current segment_start_=" << segment_start_ <<
")";
268 if (segment_start_ < 0) {
269 segment_start_ = info->start_timestamp;
270 DVLOG(2) <<
"TextChunker: Initialized segment_start_ from SegmentInfo: "
279 const int64_t kPtsWrapThreshold = 4294967296LL;
281 if (segment_end_boundary < segment_start_) {
282 int64_t diff = segment_start_ - segment_end_boundary;
283 if (diff > kPtsWrapThreshold) {
286 DVLOG(2) <<
"TextChunker: Detected PTS wrap-around. End boundary "
287 << segment_end_boundary
288 <<
" appears earlier than segment_start_ " << segment_start_
289 <<
" by " << diff <<
" ticks, but is likely "
290 <<
"later due to wrap-around.";
293 int64_t segment_end = segment_start_ + segment_duration_;
294 AddOngoingCuesToCurrentSegment(segment_end);
295 RETURN_IF_ERROR(DispatchSegment(segment_duration_));
296 segment_start_ = info->start_timestamp;
300 LOG(WARNING) <<
"TextChunker: Received SegmentInfo end boundary "
301 << segment_end_boundary <<
" that is earlier than current "
302 <<
"segment_start_ " << segment_start_ <<
" (diff: " << diff
303 <<
"). Skipping this SegmentInfo.";
309 while (segment_start_ < segment_end_boundary) {
310 int64_t segment_end = segment_start_ + segment_duration_;
314 if (segment_end > segment_end_boundary) {
315 int64_t adjusted_duration = segment_end_boundary - segment_start_;
316 if (adjusted_duration > 0) {
317 DVLOG(3) <<
"TextChunker: Dispatching adjusted segment to align with "
318 <<
"end boundary. Duration: " << adjusted_duration
319 <<
" (normal: " << segment_duration_ <<
")";
321 AddOngoingCuesToCurrentSegment(segment_end_boundary);
322 RETURN_IF_ERROR(DispatchSegment(adjusted_duration));
328 DVLOG(3) <<
"TextChunker: Dispatching full segment aligned to end boundary";
330 AddOngoingCuesToCurrentSegment(segment_end);
331 RETURN_IF_ERROR(DispatchSegment(segment_duration_));
335 segment_start_ = segment_end_boundary;
340Status TextChunker::DispatchSegment(int64_t duration) {
341 DCHECK_GT(duration, 0) <<
"Segment duration should always be positive";
343 int64_t segment_end = segment_start_ + duration;
347 DVLOG(1) <<
"DispatchSegment, start=" << segment_start_
348 <<
" end=" << segment_end;
349 for (
const auto& sample : samples_in_current_segment_) {
351 if (PtsIsBefore(sample->start_time(), segment_end)) {
352 DVLOG(2) <<
"DispatchTextSample, pts=" << sample->start_time()
353 <<
" end=" << sample->EndTime();
354 RETURN_IF_ERROR(DispatchTextSample(kStreamIndex, sample));
356 DVLOG(2) <<
"Skipping sample pts=" << sample->start_time()
357 <<
" (after segment_end=" << segment_end <<
")";
362 std::shared_ptr<SegmentInfo> info = std::make_shared<SegmentInfo>();
363 info->start_timestamp = segment_start_;
364 info->duration = duration;
365 info->segment_number = segment_number_++;
367 RETURN_IF_ERROR(DispatchSegmentInfo(kStreamIndex, std::move(info)));
370 const int64_t new_segment_start = segment_start_ + duration;
371 segment_start_ = new_segment_start;
375 samples_in_current_segment_.remove_if(
376 [new_segment_start](
const std::shared_ptr<const TextSample>& sample) {
378 return PtsIsBeforeOrEqual(sample->EndTime(), new_segment_start);
384int64_t TextChunker::ScaleTime(
double seconds)
const {
385 DCHECK_GT(time_scale_, 0) <<
"Need positive time scale to scale time.";
386 return static_cast<int64_t
>(seconds * time_scale_);
389void TextChunker::AddOngoingCuesToCurrentSegment(int64_t segment_end) {
392 for (
const auto& s : samples_without_end_) {
393 if (s->role() == TextSampleRole::kCueStart) {
395 if (PtsIsBefore(s->start_time(), segment_end)) {
397 auto cue_start = s->start_time();
398 if (PtsIsBefore(cue_start, segment_start_)) {
399 cue_start = segment_start_;
401 auto cropped_cue = std::make_shared<TextSample>(
402 "", cue_start, segment_end, s->settings(), s->body());
403 DVLOG(3) <<
"AddOngoingCuesToCurrentSegment: cropped cue start="
404 << cue_start <<
" end=" << segment_end;
405 samples_in_current_segment_.push_back(std::move(cropped_cue));
All the methods that are virtual are virtual for mocking.