Project import generated by Copybara.
GitOrigin-RevId: 53a42bf7ad836321123cb7b6c80b0f2e13fbf83e
This commit is contained in:
@@ -45,6 +45,17 @@ DefaultInputStreamHandler::DefaultInputStreamHandler(
|
||||
}
|
||||
}
|
||||
|
||||
void DefaultInputStreamHandler::PrepareForRun(
|
||||
std::function<void()> headers_ready_callback,
|
||||
std::function<void()> notification_callback,
|
||||
std::function<void(CalculatorContext*)> schedule_callback,
|
||||
std::function<void(::mediapipe::Status)> error_callback) {
|
||||
sync_set_.PrepareForRun();
|
||||
InputStreamHandler::PrepareForRun(
|
||||
std::move(headers_ready_callback), std::move(notification_callback),
|
||||
std::move(schedule_callback), std::move(error_callback));
|
||||
}
|
||||
|
||||
NodeReadiness DefaultInputStreamHandler::GetNodeReadiness(
|
||||
Timestamp* min_stream_timestamp) {
|
||||
return sync_set_.GetReadiness(min_stream_timestamp);
|
||||
|
||||
@@ -35,6 +35,13 @@ class DefaultInputStreamHandler : public InputStreamHandler {
|
||||
bool calculator_run_in_parallel);
|
||||
|
||||
protected:
|
||||
// Reinitializes this InputStreamHandler before each CalculatorGraph run.
|
||||
void PrepareForRun(
|
||||
std::function<void()> headers_ready_callback,
|
||||
std::function<void()> notification_callback,
|
||||
std::function<void(CalculatorContext*)> schedule_callback,
|
||||
std::function<void(::mediapipe::Status)> error_callback) override;
|
||||
|
||||
// In DefaultInputStreamHandler, a node is "ready" if:
|
||||
// - all streams are done (need to call Close() in this case), or
|
||||
// - the minimum bound (over all empty streams) is greater than the smallest
|
||||
|
||||
@@ -40,6 +40,13 @@ class ImmediateInputStreamHandler : public InputStreamHandler {
|
||||
const MediaPipeOptions& options, bool calculator_run_in_parallel);
|
||||
|
||||
protected:
|
||||
// Reinitializes this InputStreamHandler before each CalculatorGraph run.
|
||||
void PrepareForRun(
|
||||
std::function<void()> headers_ready_callback,
|
||||
std::function<void()> notification_callback,
|
||||
std::function<void(CalculatorContext*)> schedule_callback,
|
||||
std::function<void(::mediapipe::Status)> error_callback) override;
|
||||
|
||||
// Returns kReadyForProcess whenever a Packet is available at any of
|
||||
// the input streams, or any input stream becomes done.
|
||||
NodeReadiness GetNodeReadiness(Timestamp* min_stream_timestamp) override;
|
||||
@@ -69,6 +76,23 @@ ImmediateInputStreamHandler::ImmediateInputStreamHandler(
|
||||
}
|
||||
}
|
||||
|
||||
void ImmediateInputStreamHandler::PrepareForRun(
|
||||
std::function<void()> headers_ready_callback,
|
||||
std::function<void()> notification_callback,
|
||||
std::function<void(CalculatorContext*)> schedule_callback,
|
||||
std::function<void(::mediapipe::Status)> error_callback) {
|
||||
{
|
||||
absl::MutexLock lock(&mutex_);
|
||||
for (int i = 0; i < sync_sets_.size(); ++i) {
|
||||
sync_sets_[i].PrepareForRun();
|
||||
ready_timestamps_[i] = Timestamp::Unset();
|
||||
}
|
||||
}
|
||||
InputStreamHandler::PrepareForRun(
|
||||
std::move(headers_ready_callback), std::move(notification_callback),
|
||||
std::move(schedule_callback), std::move(error_callback));
|
||||
}
|
||||
|
||||
NodeReadiness ImmediateInputStreamHandler::GetNodeReadiness(
|
||||
Timestamp* min_stream_timestamp) {
|
||||
absl::MutexLock lock(&mutex_);
|
||||
|
||||
Reference in New Issue
Block a user