Project import generated by Copybara.
GitOrigin-RevId: bbbbcb4f5174dea33525729ede47c770069157cd
This commit is contained in:
@@ -244,6 +244,8 @@ cc_test(
|
||||
srcs = ["mux_input_stream_handler_test.cc"],
|
||||
deps = [
|
||||
":mux_input_stream_handler",
|
||||
"//mediapipe/calculators/core:gate_calculator",
|
||||
"//mediapipe/calculators/core:make_pair_calculator",
|
||||
"//mediapipe/calculators/core:mux_calculator",
|
||||
"//mediapipe/calculators/core:pass_through_calculator",
|
||||
"//mediapipe/calculators/core:round_robin_demux_calculator",
|
||||
|
||||
@@ -75,13 +75,30 @@ class MuxInputStreamHandler : public InputStreamHandler {
|
||||
int control_value = control_packet.Get<int>();
|
||||
CHECK_LE(0, control_value);
|
||||
CHECK_LT(control_value, input_stream_managers_.NumEntries() - 1);
|
||||
|
||||
const auto& data_stream = input_stream_managers_.Get(
|
||||
input_stream_managers_.BeginId() + control_value);
|
||||
|
||||
// Data stream may contain some outdated packets which failed to be popped
|
||||
// out during "FillInputSet". (This handler doesn't sync input streams,
|
||||
// hence "FillInputSet" can be triggerred before every input stream is
|
||||
// filled with packets corresponding to the same timestamp.)
|
||||
data_stream->ErasePacketsEarlierThan(*min_stream_timestamp);
|
||||
Timestamp stream_timestamp = data_stream->MinTimestampOrBound(&empty);
|
||||
if (empty) {
|
||||
CHECK_LE(stream_timestamp, *min_stream_timestamp);
|
||||
return NodeReadiness::kNotReady;
|
||||
if (stream_timestamp <= *min_stream_timestamp) {
|
||||
// "data_stream" didn't receive a packet corresponding to the current
|
||||
// "control_stream" packet yet.
|
||||
return NodeReadiness::kNotReady;
|
||||
}
|
||||
// "data_stream" timestamp bound update detected.
|
||||
return NodeReadiness::kReadyForProcess;
|
||||
}
|
||||
if (stream_timestamp > *min_stream_timestamp) {
|
||||
// The earliest packet "data_stream" holds corresponds to a control packet
|
||||
// yet to arrive, which means there won't be a "data_stream" packet
|
||||
// corresponding to the current "control_stream" packet, which should be
|
||||
// indicated as timestamp boun update.
|
||||
return NodeReadiness::kReadyForProcess;
|
||||
}
|
||||
CHECK_EQ(stream_timestamp, *min_stream_timestamp);
|
||||
return NodeReadiness::kReadyForProcess;
|
||||
|
||||
@@ -12,6 +12,7 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
#include "absl/status/status.h"
|
||||
#include "mediapipe/framework/calculator_framework.h"
|
||||
#include "mediapipe/framework/port/gmock.h"
|
||||
#include "mediapipe/framework/port/gtest.h"
|
||||
@@ -19,9 +20,10 @@
|
||||
#include "mediapipe/framework/port/status_matchers.h"
|
||||
|
||||
namespace mediapipe {
|
||||
|
||||
namespace {
|
||||
|
||||
using ::testing::ElementsAre;
|
||||
|
||||
// A regression test for b/31620439. MuxInputStreamHandler's accesses to the
|
||||
// control and data streams should be atomic so that it has a consistent view
|
||||
// of the two streams. None of the CHECKs in the GetNodeReadiness() method of
|
||||
@@ -87,5 +89,561 @@ TEST(MuxInputStreamHandlerTest, AtomicAccessToControlAndDataStreams) {
|
||||
MP_ASSERT_OK(graph.WaitUntilDone());
|
||||
}
|
||||
|
||||
MATCHER_P2(IntPacket, value, ts, "") {
|
||||
return arg.template Get<int>() == value && arg.Timestamp() == ts;
|
||||
}
|
||||
|
||||
struct GateAndMuxGraphInput {
|
||||
int input0;
|
||||
int input1;
|
||||
int input2;
|
||||
int select;
|
||||
bool allow0;
|
||||
bool allow1;
|
||||
bool allow2;
|
||||
Timestamp at;
|
||||
};
|
||||
|
||||
constexpr char kGateAndMuxGraph[] = R"pb(
|
||||
input_stream: "input0"
|
||||
input_stream: "input1"
|
||||
input_stream: "input2"
|
||||
input_stream: "select"
|
||||
input_stream: "allow0"
|
||||
input_stream: "allow1"
|
||||
input_stream: "allow2"
|
||||
node {
|
||||
calculator: "GateCalculator"
|
||||
input_stream: "ALLOW:allow0"
|
||||
input_stream: "input0"
|
||||
output_stream: "output0"
|
||||
}
|
||||
node {
|
||||
calculator: "GateCalculator"
|
||||
input_stream: "ALLOW:allow1"
|
||||
input_stream: "input1"
|
||||
output_stream: "output1"
|
||||
}
|
||||
node {
|
||||
calculator: "GateCalculator"
|
||||
input_stream: "ALLOW:allow2"
|
||||
input_stream: "input2"
|
||||
output_stream: "output2"
|
||||
}
|
||||
node {
|
||||
calculator: "MuxCalculator"
|
||||
input_stream: "INPUT:0:output0"
|
||||
input_stream: "INPUT:1:output1"
|
||||
input_stream: "INPUT:2:output2"
|
||||
input_stream: "SELECT:select"
|
||||
output_stream: "OUTPUT:output"
|
||||
input_stream_handler { input_stream_handler: "MuxInputStreamHandler" }
|
||||
})pb";
|
||||
|
||||
absl::Status SendInput(GateAndMuxGraphInput in, CalculatorGraph& graph) {
|
||||
MP_RETURN_IF_ERROR(graph.AddPacketToInputStream(
|
||||
"input0", MakePacket<int>(in.input0).At(in.at)));
|
||||
MP_RETURN_IF_ERROR(graph.AddPacketToInputStream(
|
||||
"input1", MakePacket<int>(in.input1).At(in.at)));
|
||||
MP_RETURN_IF_ERROR(graph.AddPacketToInputStream(
|
||||
"input2", MakePacket<int>(in.input2).At(in.at)));
|
||||
MP_RETURN_IF_ERROR(graph.AddPacketToInputStream(
|
||||
"select", MakePacket<int>(in.select).At(in.at)));
|
||||
MP_RETURN_IF_ERROR(graph.AddPacketToInputStream(
|
||||
"allow0", MakePacket<bool>(in.allow0).At(in.at)));
|
||||
MP_RETURN_IF_ERROR(graph.AddPacketToInputStream(
|
||||
"allow1", MakePacket<bool>(in.allow1).At(in.at)));
|
||||
MP_RETURN_IF_ERROR(graph.AddPacketToInputStream(
|
||||
"allow2", MakePacket<bool>(in.allow2).At(in.at)));
|
||||
return graph.WaitUntilIdle();
|
||||
}
|
||||
|
||||
TEST(MuxInputStreamHandlerTest, BasicMuxing) {
|
||||
CalculatorGraphConfig config =
|
||||
mediapipe::ParseTextProtoOrDie<CalculatorGraphConfig>(kGateAndMuxGraph);
|
||||
std::vector<Packet> output_packets;
|
||||
tool::AddVectorSink("output", &config, &output_packets);
|
||||
|
||||
CalculatorGraph graph;
|
||||
MP_ASSERT_OK(graph.Initialize(config));
|
||||
MP_ASSERT_OK(graph.StartRun({}));
|
||||
|
||||
MP_ASSERT_OK(SendInput({.input0 = 1000,
|
||||
.input1 = 900,
|
||||
.input2 = 800,
|
||||
.select = 0,
|
||||
.allow0 = true,
|
||||
.allow1 = false,
|
||||
.allow2 = false,
|
||||
.at = Timestamp(1)},
|
||||
graph));
|
||||
EXPECT_THAT(output_packets, ElementsAre(IntPacket(1000, Timestamp(1))));
|
||||
|
||||
MP_ASSERT_OK(SendInput({.input0 = 1000,
|
||||
.input1 = 900,
|
||||
.input2 = 800,
|
||||
.select = 1,
|
||||
.allow0 = false,
|
||||
.allow1 = true,
|
||||
.allow2 = false,
|
||||
.at = Timestamp(2)},
|
||||
graph));
|
||||
EXPECT_THAT(output_packets, ElementsAre(IntPacket(1000, Timestamp(1)),
|
||||
IntPacket(900, Timestamp(2))));
|
||||
|
||||
MP_ASSERT_OK(SendInput({.input0 = 1000,
|
||||
.input1 = 900,
|
||||
.input2 = 800,
|
||||
.select = 2,
|
||||
.allow0 = false,
|
||||
.allow1 = false,
|
||||
.allow2 = true,
|
||||
.at = Timestamp(3)},
|
||||
graph));
|
||||
EXPECT_THAT(output_packets, ElementsAre(IntPacket(1000, Timestamp(1)),
|
||||
IntPacket(900, Timestamp(2)),
|
||||
IntPacket(800, Timestamp(3))));
|
||||
|
||||
MP_ASSERT_OK(graph.CloseAllInputStreams());
|
||||
MP_ASSERT_OK(graph.WaitUntilDone());
|
||||
}
|
||||
|
||||
TEST(MuxInputStreamHandlerTest, MuxingNonEmptyInputs) {
|
||||
CalculatorGraphConfig config =
|
||||
mediapipe::ParseTextProtoOrDie<CalculatorGraphConfig>(kGateAndMuxGraph);
|
||||
std::vector<Packet> output_packets;
|
||||
tool::AddVectorSink("output", &config, &output_packets);
|
||||
|
||||
CalculatorGraph graph;
|
||||
MP_ASSERT_OK(graph.Initialize(config));
|
||||
MP_ASSERT_OK(graph.StartRun({}));
|
||||
|
||||
MP_ASSERT_OK(SendInput({.input0 = 1000,
|
||||
.input1 = 900,
|
||||
.input2 = 800,
|
||||
.select = 0,
|
||||
.allow0 = true,
|
||||
.allow1 = true,
|
||||
.allow2 = true,
|
||||
.at = Timestamp(1)},
|
||||
graph));
|
||||
EXPECT_THAT(output_packets, ElementsAre(IntPacket(1000, Timestamp(1))));
|
||||
|
||||
MP_ASSERT_OK(SendInput({.input0 = 1000,
|
||||
.input1 = 900,
|
||||
.input2 = 800,
|
||||
.select = 1,
|
||||
.allow0 = true,
|
||||
.allow1 = true,
|
||||
.allow2 = true,
|
||||
.at = Timestamp(2)},
|
||||
graph));
|
||||
EXPECT_THAT(output_packets, ElementsAre(IntPacket(1000, Timestamp(1)),
|
||||
IntPacket(900, Timestamp(2))));
|
||||
|
||||
MP_ASSERT_OK(SendInput({.input0 = 1000,
|
||||
.input1 = 900,
|
||||
.input2 = 800,
|
||||
.select = 2,
|
||||
.allow0 = true,
|
||||
.allow1 = true,
|
||||
.allow2 = true,
|
||||
.at = Timestamp(3)},
|
||||
graph));
|
||||
EXPECT_THAT(output_packets, ElementsAre(IntPacket(1000, Timestamp(1)),
|
||||
IntPacket(900, Timestamp(2)),
|
||||
IntPacket(800, Timestamp(3))));
|
||||
|
||||
MP_ASSERT_OK(graph.CloseAllInputStreams());
|
||||
MP_ASSERT_OK(graph.WaitUntilDone());
|
||||
}
|
||||
|
||||
TEST(MuxInputStreamHandlerTest, MuxingAllTimestampBoundUpdates) {
|
||||
CalculatorGraphConfig config =
|
||||
mediapipe::ParseTextProtoOrDie<CalculatorGraphConfig>(kGateAndMuxGraph);
|
||||
std::vector<Packet> output_packets;
|
||||
tool::AddVectorSink("output", &config, &output_packets);
|
||||
|
||||
CalculatorGraph graph;
|
||||
MP_ASSERT_OK(graph.Initialize(config));
|
||||
MP_ASSERT_OK(graph.StartRun({}));
|
||||
|
||||
MP_ASSERT_OK(SendInput({.input0 = 1000,
|
||||
.input1 = 900,
|
||||
.input2 = 800,
|
||||
.select = 0,
|
||||
.allow0 = false,
|
||||
.allow1 = false,
|
||||
.allow2 = false,
|
||||
.at = Timestamp(1)},
|
||||
graph));
|
||||
EXPECT_TRUE(output_packets.empty());
|
||||
|
||||
MP_ASSERT_OK(SendInput({.input0 = 1000,
|
||||
.input1 = 900,
|
||||
.input2 = 800,
|
||||
.select = 1,
|
||||
.allow0 = false,
|
||||
.allow1 = false,
|
||||
.allow2 = false,
|
||||
.at = Timestamp(2)},
|
||||
graph));
|
||||
EXPECT_TRUE(output_packets.empty());
|
||||
|
||||
MP_ASSERT_OK(SendInput({.input0 = 1000,
|
||||
.input1 = 900,
|
||||
.input2 = 800,
|
||||
.select = 2,
|
||||
.allow0 = false,
|
||||
.allow1 = false,
|
||||
.allow2 = false,
|
||||
.at = Timestamp(3)},
|
||||
graph));
|
||||
EXPECT_TRUE(output_packets.empty());
|
||||
|
||||
MP_ASSERT_OK(graph.CloseAllInputStreams());
|
||||
MP_ASSERT_OK(graph.WaitUntilDone());
|
||||
}
|
||||
|
||||
TEST(MuxInputStreamHandlerTest, MuxingSlectedTimestampBoundUpdates) {
|
||||
CalculatorGraphConfig config =
|
||||
mediapipe::ParseTextProtoOrDie<CalculatorGraphConfig>(kGateAndMuxGraph);
|
||||
std::vector<Packet> output_packets;
|
||||
tool::AddVectorSink("output", &config, &output_packets);
|
||||
|
||||
CalculatorGraph graph;
|
||||
MP_ASSERT_OK(graph.Initialize(config));
|
||||
MP_ASSERT_OK(graph.StartRun({}));
|
||||
|
||||
MP_ASSERT_OK(SendInput({.input0 = 1000,
|
||||
.input1 = 900,
|
||||
.input2 = 800,
|
||||
.select = 0,
|
||||
.allow0 = false,
|
||||
.allow1 = true,
|
||||
.allow2 = true,
|
||||
.at = Timestamp(1)},
|
||||
graph));
|
||||
EXPECT_TRUE(output_packets.empty());
|
||||
|
||||
MP_ASSERT_OK(SendInput({.input0 = 1000,
|
||||
.input1 = 900,
|
||||
.input2 = 800,
|
||||
.select = 1,
|
||||
.allow0 = true,
|
||||
.allow1 = false,
|
||||
.allow2 = true,
|
||||
.at = Timestamp(2)},
|
||||
graph));
|
||||
EXPECT_TRUE(output_packets.empty());
|
||||
|
||||
MP_ASSERT_OK(SendInput({.input0 = 1000,
|
||||
.input1 = 900,
|
||||
.input2 = 800,
|
||||
.select = 2,
|
||||
.allow0 = true,
|
||||
.allow1 = true,
|
||||
.allow2 = false,
|
||||
.at = Timestamp(3)},
|
||||
graph));
|
||||
EXPECT_TRUE(output_packets.empty());
|
||||
|
||||
MP_ASSERT_OK(graph.CloseAllInputStreams());
|
||||
MP_ASSERT_OK(graph.WaitUntilDone());
|
||||
}
|
||||
|
||||
TEST(MuxInputStreamHandlerTest, MuxingSometimesTimestampBoundUpdates) {
|
||||
CalculatorGraphConfig config =
|
||||
mediapipe::ParseTextProtoOrDie<CalculatorGraphConfig>(kGateAndMuxGraph);
|
||||
std::vector<Packet> output_packets;
|
||||
tool::AddVectorSink("output", &config, &output_packets);
|
||||
|
||||
CalculatorGraph graph;
|
||||
MP_ASSERT_OK(graph.Initialize(config));
|
||||
MP_ASSERT_OK(graph.StartRun({}));
|
||||
|
||||
MP_ASSERT_OK(SendInput({.input0 = 1000,
|
||||
.input1 = 900,
|
||||
.input2 = 800,
|
||||
.select = 0,
|
||||
.allow0 = false,
|
||||
.allow1 = false,
|
||||
.allow2 = false,
|
||||
.at = Timestamp(1)},
|
||||
graph));
|
||||
EXPECT_TRUE(output_packets.empty());
|
||||
|
||||
MP_ASSERT_OK(SendInput({.input0 = 1000,
|
||||
.input1 = 900,
|
||||
.input2 = 800,
|
||||
.select = 1,
|
||||
.allow0 = false,
|
||||
.allow1 = true,
|
||||
.allow2 = false,
|
||||
.at = Timestamp(2)},
|
||||
graph));
|
||||
EXPECT_THAT(output_packets, ElementsAre(IntPacket(900, Timestamp(2))));
|
||||
|
||||
MP_ASSERT_OK(SendInput({.input0 = 1000,
|
||||
.input1 = 900,
|
||||
.input2 = 800,
|
||||
.select = 2,
|
||||
.allow0 = true,
|
||||
.allow1 = true,
|
||||
.allow2 = false,
|
||||
.at = Timestamp(3)},
|
||||
graph));
|
||||
EXPECT_THAT(output_packets, ElementsAre(IntPacket(900, Timestamp(2))));
|
||||
|
||||
MP_ASSERT_OK(SendInput({.input0 = 1000,
|
||||
.input1 = 900,
|
||||
.input2 = 800,
|
||||
.select = 2,
|
||||
.allow0 = true,
|
||||
.allow1 = true,
|
||||
.allow2 = true,
|
||||
.at = Timestamp(4)},
|
||||
graph));
|
||||
EXPECT_THAT(output_packets, ElementsAre(IntPacket(900, Timestamp(2)),
|
||||
IntPacket(800, Timestamp(4))));
|
||||
|
||||
MP_ASSERT_OK(SendInput({.input0 = 700,
|
||||
.input1 = 600,
|
||||
.input2 = 500,
|
||||
.select = 0,
|
||||
.allow0 = true,
|
||||
.allow1 = false,
|
||||
.allow2 = false,
|
||||
.at = Timestamp(5)},
|
||||
graph));
|
||||
EXPECT_THAT(output_packets, ElementsAre(IntPacket(900, Timestamp(2)),
|
||||
IntPacket(800, Timestamp(4)),
|
||||
IntPacket(700, Timestamp(5))));
|
||||
|
||||
MP_ASSERT_OK(graph.CloseAllInputStreams());
|
||||
MP_ASSERT_OK(graph.WaitUntilDone());
|
||||
}
|
||||
|
||||
MATCHER_P(EmptyPacket, ts, "") {
|
||||
return arg.IsEmpty() && arg.Timestamp() == ts;
|
||||
}
|
||||
|
||||
MATCHER_P2(Pair, m1, m2, "") {
|
||||
const auto& p = arg.template Get<std::pair<Packet, Packet>>();
|
||||
return testing::Matches(m1)(p.first) && testing::Matches(m2)(p.second);
|
||||
}
|
||||
|
||||
TEST(MuxInputStreamHandlerTest,
|
||||
TimestampBoundUpdateWhenControlPacketEarlierThanDataPacket) {
|
||||
CalculatorGraphConfig config =
|
||||
mediapipe::ParseTextProtoOrDie<CalculatorGraphConfig>(R"pb(
|
||||
input_stream: "input0"
|
||||
input_stream: "input1"
|
||||
input_stream: "select"
|
||||
node {
|
||||
calculator: "MuxCalculator"
|
||||
input_stream: "INPUT:0:input0"
|
||||
input_stream: "INPUT:1:input1"
|
||||
input_stream: "SELECT:select"
|
||||
output_stream: "OUTPUT:output"
|
||||
input_stream_handler { input_stream_handler: "MuxInputStreamHandler" }
|
||||
}
|
||||
node {
|
||||
calculator: "MakePairCalculator"
|
||||
input_stream: "select"
|
||||
input_stream: "output"
|
||||
output_stream: "pair"
|
||||
}
|
||||
)pb");
|
||||
std::vector<Packet> pair_packets;
|
||||
tool::AddVectorSink("pair", &config, &pair_packets);
|
||||
|
||||
CalculatorGraph graph;
|
||||
MP_ASSERT_OK(graph.Initialize(config));
|
||||
MP_ASSERT_OK(graph.StartRun({}));
|
||||
|
||||
MP_ASSERT_OK(graph.AddPacketToInputStream(
|
||||
"select", MakePacket<int>(0).At(Timestamp(1))));
|
||||
MP_ASSERT_OK(graph.WaitUntilIdle());
|
||||
EXPECT_TRUE(pair_packets.empty());
|
||||
|
||||
MP_ASSERT_OK(graph.AddPacketToInputStream(
|
||||
"input0", MakePacket<int>(1000).At(Timestamp(2))));
|
||||
|
||||
MP_ASSERT_OK(graph.WaitUntilIdle());
|
||||
EXPECT_THAT(pair_packets, ElementsAre(Pair(IntPacket(0, Timestamp(1)),
|
||||
EmptyPacket(Timestamp(1)))));
|
||||
|
||||
MP_ASSERT_OK(graph.AddPacketToInputStream(
|
||||
"input1", MakePacket<int>(900).At(Timestamp(1))));
|
||||
MP_ASSERT_OK(graph.AddPacketToInputStream(
|
||||
"input1", MakePacket<int>(800).At(Timestamp(4))));
|
||||
|
||||
MP_ASSERT_OK(graph.WaitUntilIdle());
|
||||
EXPECT_THAT(pair_packets, ElementsAre(Pair(IntPacket(0, Timestamp(1)),
|
||||
EmptyPacket(Timestamp(1)))));
|
||||
|
||||
MP_ASSERT_OK(graph.AddPacketToInputStream(
|
||||
"select", MakePacket<int>(0).At(Timestamp(2))));
|
||||
MP_ASSERT_OK(graph.WaitUntilIdle());
|
||||
EXPECT_THAT(
|
||||
pair_packets,
|
||||
ElementsAre(
|
||||
Pair(IntPacket(0, Timestamp(1)), EmptyPacket(Timestamp(1))),
|
||||
Pair(IntPacket(0, Timestamp(2)), IntPacket(1000, Timestamp(2)))));
|
||||
|
||||
MP_ASSERT_OK(graph.AddPacketToInputStream(
|
||||
"select", MakePacket<int>(1).At(Timestamp(3))));
|
||||
MP_ASSERT_OK(graph.WaitUntilIdle());
|
||||
EXPECT_THAT(
|
||||
pair_packets,
|
||||
ElementsAre(
|
||||
Pair(IntPacket(0, Timestamp(1)), EmptyPacket(Timestamp(1))),
|
||||
Pair(IntPacket(0, Timestamp(2)), IntPacket(1000, Timestamp(2))),
|
||||
Pair(IntPacket(1, Timestamp(3)), EmptyPacket(Timestamp(3)))));
|
||||
|
||||
MP_ASSERT_OK(graph.AddPacketToInputStream(
|
||||
"select", MakePacket<int>(1).At(Timestamp(4))));
|
||||
MP_ASSERT_OK(graph.WaitUntilIdle());
|
||||
EXPECT_THAT(
|
||||
pair_packets,
|
||||
ElementsAre(
|
||||
Pair(IntPacket(0, Timestamp(1)), EmptyPacket(Timestamp(1))),
|
||||
Pair(IntPacket(0, Timestamp(2)), IntPacket(1000, Timestamp(2))),
|
||||
Pair(IntPacket(1, Timestamp(3)), EmptyPacket(Timestamp(3))),
|
||||
Pair(IntPacket(1, Timestamp(4)), IntPacket(800, Timestamp(4)))));
|
||||
|
||||
MP_ASSERT_OK(graph.CloseAllInputStreams());
|
||||
MP_ASSERT_OK(graph.WaitUntilDone());
|
||||
}
|
||||
|
||||
TEST(MuxInputStreamHandlerTest,
|
||||
TimestampBoundUpdateWhenControlPacketEarlierThanDataPacketPacketsAtOnce) {
|
||||
CalculatorGraphConfig config =
|
||||
mediapipe::ParseTextProtoOrDie<CalculatorGraphConfig>(R"pb(
|
||||
input_stream: "input0"
|
||||
input_stream: "input1"
|
||||
input_stream: "select"
|
||||
node {
|
||||
calculator: "MuxCalculator"
|
||||
input_stream: "INPUT:0:input0"
|
||||
input_stream: "INPUT:1:input1"
|
||||
input_stream: "SELECT:select"
|
||||
output_stream: "OUTPUT:output"
|
||||
input_stream_handler { input_stream_handler: "MuxInputStreamHandler" }
|
||||
}
|
||||
node {
|
||||
calculator: "MakePairCalculator"
|
||||
input_stream: "select"
|
||||
input_stream: "output"
|
||||
output_stream: "pair"
|
||||
}
|
||||
)pb");
|
||||
std::vector<Packet> pair_packets;
|
||||
tool::AddVectorSink("pair", &config, &pair_packets);
|
||||
|
||||
CalculatorGraph graph;
|
||||
MP_ASSERT_OK(graph.Initialize(config));
|
||||
MP_ASSERT_OK(graph.StartRun({}));
|
||||
|
||||
MP_ASSERT_OK(graph.AddPacketToInputStream(
|
||||
"select", MakePacket<int>(0).At(Timestamp(1))));
|
||||
MP_ASSERT_OK(graph.AddPacketToInputStream(
|
||||
"input0", MakePacket<int>(1000).At(Timestamp(2))));
|
||||
MP_ASSERT_OK(graph.AddPacketToInputStream(
|
||||
"input1", MakePacket<int>(900).At(Timestamp(1))));
|
||||
MP_ASSERT_OK(graph.AddPacketToInputStream(
|
||||
"input1", MakePacket<int>(800).At(Timestamp(4))));
|
||||
MP_ASSERT_OK(graph.AddPacketToInputStream(
|
||||
"select", MakePacket<int>(0).At(Timestamp(2))));
|
||||
MP_ASSERT_OK(graph.AddPacketToInputStream(
|
||||
"select", MakePacket<int>(1).At(Timestamp(3))));
|
||||
MP_ASSERT_OK(graph.AddPacketToInputStream(
|
||||
"select", MakePacket<int>(1).At(Timestamp(4))));
|
||||
MP_ASSERT_OK(graph.WaitUntilIdle());
|
||||
EXPECT_THAT(
|
||||
pair_packets,
|
||||
ElementsAre(
|
||||
Pair(IntPacket(0, Timestamp(1)), EmptyPacket(Timestamp(1))),
|
||||
Pair(IntPacket(0, Timestamp(2)), IntPacket(1000, Timestamp(2))),
|
||||
Pair(IntPacket(1, Timestamp(3)), EmptyPacket(Timestamp(3))),
|
||||
Pair(IntPacket(1, Timestamp(4)), IntPacket(800, Timestamp(4)))));
|
||||
|
||||
MP_ASSERT_OK(graph.CloseAllInputStreams());
|
||||
MP_ASSERT_OK(graph.WaitUntilDone());
|
||||
}
|
||||
|
||||
TEST(MuxInputStreamHandlerTest,
|
||||
TimestampBoundUpdateTriggersTimestampBoundUpdate) {
|
||||
CalculatorGraphConfig config =
|
||||
mediapipe::ParseTextProtoOrDie<CalculatorGraphConfig>(R"pb(
|
||||
input_stream: "input0"
|
||||
input_stream: "input1"
|
||||
input_stream: "select"
|
||||
input_stream: "allow0"
|
||||
input_stream: "allow1"
|
||||
node {
|
||||
calculator: "GateCalculator"
|
||||
input_stream: "ALLOW:allow0"
|
||||
input_stream: "input0"
|
||||
output_stream: "output0"
|
||||
}
|
||||
node {
|
||||
calculator: "GateCalculator"
|
||||
input_stream: "ALLOW:allow1"
|
||||
input_stream: "input1"
|
||||
output_stream: "output1"
|
||||
}
|
||||
node {
|
||||
calculator: "MuxCalculator"
|
||||
input_stream: "INPUT:0:output0"
|
||||
input_stream: "INPUT:1:output1"
|
||||
input_stream: "SELECT:select"
|
||||
output_stream: "OUTPUT:output"
|
||||
input_stream_handler { input_stream_handler: "MuxInputStreamHandler" }
|
||||
}
|
||||
node {
|
||||
calculator: "MakePairCalculator"
|
||||
input_stream: "select"
|
||||
input_stream: "output"
|
||||
output_stream: "pair"
|
||||
}
|
||||
)pb");
|
||||
std::vector<Packet> pair_packets;
|
||||
tool::AddVectorSink("pair", &config, &pair_packets);
|
||||
|
||||
CalculatorGraph graph;
|
||||
MP_ASSERT_OK(graph.Initialize(config));
|
||||
MP_ASSERT_OK(graph.StartRun({}));
|
||||
|
||||
MP_ASSERT_OK(graph.AddPacketToInputStream(
|
||||
"select", MakePacket<int>(0).At(Timestamp(1))));
|
||||
MP_ASSERT_OK(graph.AddPacketToInputStream(
|
||||
"input0", MakePacket<int>(1000).At(Timestamp(1))));
|
||||
MP_ASSERT_OK(graph.AddPacketToInputStream(
|
||||
"allow0", MakePacket<bool>(false).At(Timestamp(1))));
|
||||
|
||||
MP_ASSERT_OK(graph.WaitUntilIdle());
|
||||
EXPECT_THAT(pair_packets, ElementsAre(Pair(IntPacket(0, Timestamp(1)),
|
||||
EmptyPacket(Timestamp(1)))));
|
||||
|
||||
MP_ASSERT_OK(graph.AddPacketToInputStream(
|
||||
"select", MakePacket<int>(0).At(Timestamp(2))));
|
||||
MP_ASSERT_OK(graph.AddPacketToInputStream(
|
||||
"input0", MakePacket<int>(900).At(Timestamp(2))));
|
||||
MP_ASSERT_OK(graph.AddPacketToInputStream(
|
||||
"allow0", MakePacket<bool>(true).At(Timestamp(2))));
|
||||
|
||||
MP_ASSERT_OK(graph.WaitUntilIdle());
|
||||
EXPECT_THAT(
|
||||
pair_packets,
|
||||
ElementsAre(
|
||||
Pair(IntPacket(0, Timestamp(1)), EmptyPacket(Timestamp(1))),
|
||||
Pair(IntPacket(0, Timestamp(2)), IntPacket(900, Timestamp(2)))));
|
||||
|
||||
MP_ASSERT_OK(graph.CloseAllInputStreams());
|
||||
MP_ASSERT_OK(graph.WaitUntilDone());
|
||||
}
|
||||
|
||||
} // namespace
|
||||
} // namespace mediapipe
|
||||
|
||||
Reference in New Issue
Block a user