From 2f12ec2c79dfdf10796822b9cffa2680392911ef Mon Sep 17 00:00:00 2001 From: Marc Paterno Date: Fri, 24 Jul 2026 12:23:23 -0500 Subject: [PATCH 1/5] test: add unit test for transform_node --- test/CMakeLists.txt | 8 ++ test/transform_node_test.cpp | 227 +++++++++++++++++++++++++++++++++++ 2 files changed, 235 insertions(+) create mode 100644 test/transform_node_test.cpp diff --git a/test/CMakeLists.txt b/test/CMakeLists.txt index 245a5c8c8..7d417381e 100644 --- a/test/CMakeLists.txt +++ b/test/CMakeLists.txt @@ -209,6 +209,14 @@ cet_test( LIBRARIES phlex::core_internal ) +cet_test( + transform_node + USE_CATCH2_MAIN + SOURCE + transform_node_test.cpp + LIBRARIES + phlex::core_internal +) cet_test( replicated USE_CATCH2_MAIN diff --git a/test/transform_node_test.cpp b/test/transform_node_test.cpp new file mode 100644 index 000000000..53c557c60 --- /dev/null +++ b/test/transform_node_test.cpp @@ -0,0 +1,227 @@ +#include "phlex/core/declared_transform.hpp" +#include "phlex/metaprogramming/delegate.hpp" +#include "phlex/model/data_cell_index.hpp" +#include "phlex/model/product_store.hpp" + +#include "catch2/catch_test_macros.hpp" +#include "oneapi/tbb/flow_graph.h" + +#include +#include +#include +#include +#include + +using namespace phlex; +using namespace phlex::detail; +using namespace phlex::experimental; + +namespace { + constexpr auto message_id = 42u; + + struct input_type_1 { + int value; + auto operator<=>(input_type_1 const&) const = default; + }; + + struct input_type_2 { + int value; + auto operator<=>(input_type_2 const&) const = default; + }; + + struct output_type_1 { + int value; + auto operator<=>(output_type_1 const&) const = default; + }; + + struct output_type_2 { + std::string value; + auto operator<=>(output_type_2 const&) const = default; + }; + + output_type_1 double_value(input_type_1 const& input) { return {input.value * 2}; } + + output_type_1 add_values(input_type_1 const& left, input_type_2 const& right) + { + return {left.value + right.value}; + } + + auto number_and_label(input_type_1 const& input) + { + return std::tuple{output_type_1{input.value}, output_type_2{std::to_string(input.value)}}; + } + + template + product_specification spec(char const* creator, char const* suffix) + { + return {algorithm_name{creator}, identifier{suffix}, make_type_id()}; + } + + template + product_selector selector(char const* creator, char const* suffix) + { + return {.creator = creator, .layer = "job", .suffix = suffix, .type = make_type_id()}; + } + + template + product_store_ptr store_with_product(char const* creator, char const* suffix, T value) + { + auto store = product_store::base(creator); + store->add_product(spec(creator, suffix), std::move(value)); + return store; + } + + template + auto algorithm_bits_for(F function) + { + return algorithm_bits{std::shared_ptr{}, std::move(function)}; + } +} + +TEST_CASE("transform_node directly transforms one input product", "[transform_node]") +{ + oneapi::tbb::flow::graph graph; + auto input_selector = selector("input", ""); + auto input_store = store_with_product("input", "", input_type_1{21}); + auto alg = algorithm_bits_for(double_value); + + transform_node node{ + algorithm_name{"double_value"}, 1u, {}, graph, std::move(alg), {input_selector}, {}}; + declared_transform& transform = node; + + auto const& output_specs = transform.output(); + REQUIRE(output_specs.size() == 1u); + CHECK(output_specs[0].creator() == algorithm_name{"double_value"}); + CHECK(output_specs[0].suffix() == identifier{""}); + CHECK(output_specs[0].type() == make_type_id()); + + oneapi::tbb::flow::queue_node sink{graph}; + make_edge(transform.output_port(), sink); + + REQUIRE(node.port(input_selector).try_put({.store = input_store, .id = message_id})); + graph.wait_for_all(); + + message output; + REQUIRE(sink.try_get(output)); + + CHECK(output.id == message_id); + REQUIRE(output.store); + CHECK(output.store->index() == input_store->index()); + CHECK(output.store->source() == algorithm_name{"double_value"}); + CHECK(output.store->get_product(output_specs[0]) == output_type_1{42}); + + CHECK(transform.num_calls() == 1u); + CHECK(transform.product_count() == 1u); +} + +TEST_CASE("transform_node joins multiple input products", "[transform_node]") +{ + oneapi::tbb::flow::graph graph; + auto left_selector = selector("left_input", ""); + auto right_selector = selector("right_input", ""); + auto left_store = store_with_product("left_input", "", input_type_1{17}); + auto right_store = store_with_product("right_input", "", input_type_2{25}); + auto alg = algorithm_bits_for(add_values); + + transform_node node{algorithm_name{"add_values"}, + 1u, + {}, + graph, + std::move(alg), + {left_selector, right_selector}, + {}}; + declared_transform& transform = node; + + auto const& output_specs = transform.output(); + REQUIRE(output_specs.size() == 1u); + CHECK(output_specs[0].suffix() == identifier{""}); + + oneapi::tbb::flow::queue_node sink{graph}; + make_edge(transform.output_port(), sink); + + auto& left_port = *transform.ports().at(0); + auto& right_port = *transform.ports().at(1); + + REQUIRE(left_port.try_put({.store = left_store, .id = message_id})); + graph.wait_for_all(); + + message output; + CHECK_FALSE(sink.try_get(output)); + CHECK(transform.num_calls() == 0u); + + REQUIRE(right_port.try_put({.store = right_store, .id = message_id})); + graph.wait_for_all(); + CHECK_FALSE(sink.try_get(output)); + CHECK(transform.num_calls() == 0u); + + auto index_ports = transform.index_ports(); + REQUIRE(index_ports.size() == 2u); + REQUIRE(index_ports[0].index_port->try_put( + {.index = left_store->index(), .msg_id = message_id, .cache = true})); + graph.wait_for_all(); + CHECK_FALSE(sink.try_get(output)); + CHECK(transform.num_calls() == 0u); + + REQUIRE(index_ports[1].index_port->try_put( + {.index = right_store->index(), .msg_id = message_id, .cache = true})); + graph.wait_for_all(); + + REQUIRE(sink.try_get(output)); + CHECK_FALSE(sink.try_get(output)); + + REQUIRE(index_ports[0].token_port->try_put({.index = left_store->index(), .count = 1})); + REQUIRE(index_ports[1].token_port->try_put({.index = right_store->index(), .count = 1})); + graph.wait_for_all(); + + CHECK(output.id == message_id); + REQUIRE(output.store); + CHECK(output.store->index() == left_store->index()); + + CHECK(output.store->get_product(output_specs[0]) == output_type_1{42}); + + CHECK(transform.num_calls() == 1u); + CHECK(transform.product_count() == 1u); +} + +TEST_CASE("transform_node stores multiple output products", "[transform_node]") +{ + oneapi::tbb::flow::graph graph; + auto input_selector = selector("input", ""); + auto input_store = store_with_product("input", "", input_type_1{7}); + auto alg = algorithm_bits_for(number_and_label); + + transform_node node{algorithm_name{"number_and_label"}, + 1u, + {}, + graph, + std::move(alg), + {input_selector}, + {"number", "label"}}; + declared_transform& transform = node; + + oneapi::tbb::flow::queue_node sink{graph}; + make_edge(transform.output_port(), sink); + + REQUIRE(node.port(input_selector).try_put({.store = input_store, .id = message_id})); + graph.wait_for_all(); + + message output; + REQUIRE(sink.try_get(output)); + CHECK_FALSE(sink.try_get(output)); + + auto const& output_specs = transform.output(); + REQUIRE(output_specs.size() == 2u); + CHECK(output_specs[0].creator() == algorithm_name{"number_and_label"}); + CHECK(output_specs[0].suffix() == identifier{"number"}); + CHECK(output_specs[0].type() == make_type_id()); + CHECK(output_specs[1].creator() == algorithm_name{"number_and_label"}); + CHECK(output_specs[1].suffix() == identifier{"label"}); + CHECK(output_specs[1].type() == make_type_id()); + + REQUIRE(output.store); + CHECK(output.store->get_product(output_specs[0]) == output_type_1{7}); + CHECK(output.store->get_product(output_specs[1]) == output_type_2{"7"}); + + CHECK(transform.num_calls() == 1u); + CHECK(transform.product_count() == 1u); +} From 7e9c6e65bb2f1fca55910cc37815ff1b93c2df9b Mon Sep 17 00:00:00 2001 From: Marc Paterno Date: Tue, 28 Jul 2026 13:55:11 -0500 Subject: [PATCH 2/5] feat(test): add join test for transform_node with multiple inputs --- test/CMakeLists.txt | 8 +++ test/join_test.cpp | 135 ++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 143 insertions(+) create mode 100644 test/join_test.cpp diff --git a/test/CMakeLists.txt b/test/CMakeLists.txt index 7d417381e..1168f5d5a 100644 --- a/test/CMakeLists.txt +++ b/test/CMakeLists.txt @@ -217,6 +217,14 @@ cet_test( LIBRARIES phlex::core_internal ) +cet_test( + join + USE_CATCH2_MAIN + SOURCE + join_test.cpp + LIBRARIES + phlex::core_internal +) cet_test( replicated USE_CATCH2_MAIN diff --git a/test/join_test.cpp b/test/join_test.cpp new file mode 100644 index 000000000..dc96db681 --- /dev/null +++ b/test/join_test.cpp @@ -0,0 +1,135 @@ +#include "phlex/core/declared_transform.hpp" +#include "phlex/metaprogramming/delegate.hpp" +#include "phlex/model/data_cell_index.hpp" +#include "phlex/model/product_store.hpp" + +#include "catch2/catch_test_macros.hpp" +#include "oneapi/tbb/flow_graph.h" + +#include +#include +#include +#include +#include + +using namespace phlex; +using namespace phlex::detail; +using namespace phlex::experimental; + +namespace { + constexpr auto message_id = 42u; + + struct input_type_1 { + int value; + auto operator<=>(input_type_1 const&) const = default; + }; + + struct input_type_2 { + int value; + auto operator<=>(input_type_2 const&) const = default; + }; + + struct output_type_1 { + int value; + auto operator<=>(output_type_1 const&) const = default; + }; + + output_type_1 add_values(input_type_1 const& left, input_type_2 const& right) + { + return {left.value + right.value}; + } + + template + product_specification spec(char const* creator, char const* suffix) + { + return {algorithm_name{creator}, identifier{suffix}, make_type_id()}; + } + + template + product_selector selector(char const* creator, char const* suffix) + { + return {.creator = creator, .layer = "job", .suffix = suffix, .type = make_type_id()}; + } + + template + product_store_ptr store_with_product(char const* creator, char const* suffix, T value) + { + auto store = product_store::base(creator); + store->add_product(spec(creator, suffix), std::move(value)); + return store; + } + + template + auto algorithm_bits_for(F function) + { + return algorithm_bits{std::shared_ptr{}, std::move(function)}; + } +} + +TEST_CASE("transform_node joins multiple input products", "[join]") +{ + oneapi::tbb::flow::graph graph; + auto left_selector = selector("left_input", ""); + auto right_selector = selector("right_input", ""); + auto left_store = store_with_product("left_input", "", input_type_1{17}); + auto right_store = store_with_product("right_input", "", input_type_2{25}); + auto alg = algorithm_bits_for(add_values); + + transform_node node{algorithm_name{"add_values"}, + 1u, + {}, + graph, + std::move(alg), + {left_selector, right_selector}, + {}}; + declared_transform& transform = node; + + oneapi::tbb::flow::queue_node sink{graph}; + make_edge(transform.output_port(), sink); + + auto& left_port = *transform.ports().at(0); + auto& right_port = *transform.ports().at(1); + + REQUIRE(left_port.try_put({.store = left_store, .id = message_id})); + graph.wait_for_all(); + + message output; + CHECK_FALSE(sink.try_get(output)); + CHECK(transform.num_calls() == 0u); + + REQUIRE(right_port.try_put({.store = right_store, .id = message_id})); + graph.wait_for_all(); + CHECK_FALSE(sink.try_get(output)); + CHECK(transform.num_calls() == 0u); + + auto index_ports = transform.index_ports(); + REQUIRE(index_ports.size() == 2u); + REQUIRE(index_ports[0].index_port->try_put( + {.index = left_store->index(), .msg_id = message_id, .cache = true})); + graph.wait_for_all(); + CHECK_FALSE(sink.try_get(output)); + CHECK(transform.num_calls() == 0u); + + REQUIRE(index_ports[1].index_port->try_put( + {.index = right_store->index(), .msg_id = message_id, .cache = true})); + graph.wait_for_all(); + + REQUIRE(sink.try_get(output)); + CHECK_FALSE(sink.try_get(output)); + + REQUIRE(index_ports[0].token_port->try_put({.index = left_store->index(), .count = 1})); + REQUIRE(index_ports[1].token_port->try_put({.index = right_store->index(), .count = 1})); + graph.wait_for_all(); + + CHECK(output.id == message_id); + REQUIRE(output.store); + CHECK(output.store->index() == left_store->index()); + + auto const& output_specs = transform.output(); + REQUIRE(output_specs.size() == 1u); + CHECK(output_specs[0].suffix() == identifier{""}); + CHECK(output.store->get_product(output_specs[0]) == output_type_1{42}); + + CHECK(transform.num_calls() == 1u); + CHECK(transform.product_count() == 1u); +} From 928fe36862fb9ca34ab463a2fb49cc315d133c55 Mon Sep 17 00:00:00 2001 From: Marc Paterno Date: Fri, 31 Jul 2026 11:08:06 -0500 Subject: [PATCH 3/5] test: update join tests for multilayer_join_node Move the multi-input join coverage from `transform_node` tests to a dedicated `multilayer_join_node` test by renaming `join_test.cpp` to `multilayer_join_test.cpp`, updating the test case name, and adjusting the test wiring. Remove the old join test from `test/transform_node_test.cpp` and update `test/CMakeLists.txt`. --- test/CMakeLists.txt | 4 +- ...join_test.cpp => multilayer_join_test.cpp} | 78 ++++++------------- test/transform_node_test.cpp | 74 ------------------ 3 files changed, 26 insertions(+), 130 deletions(-) rename test/{join_test.cpp => multilayer_join_test.cpp} (51%) diff --git a/test/CMakeLists.txt b/test/CMakeLists.txt index 1168f5d5a..6f57713ca 100644 --- a/test/CMakeLists.txt +++ b/test/CMakeLists.txt @@ -218,10 +218,10 @@ cet_test( phlex::core_internal ) cet_test( - join + multilayer_join USE_CATCH2_MAIN SOURCE - join_test.cpp + multilayer_join_test.cpp LIBRARIES phlex::core_internal ) diff --git a/test/join_test.cpp b/test/multilayer_join_test.cpp similarity index 51% rename from test/join_test.cpp rename to test/multilayer_join_test.cpp index dc96db681..1c04fd738 100644 --- a/test/join_test.cpp +++ b/test/multilayer_join_test.cpp @@ -1,5 +1,4 @@ -#include "phlex/core/declared_transform.hpp" -#include "phlex/metaprogramming/delegate.hpp" +#include "phlex/core/multilayer_join_node.hpp" #include "phlex/model/data_cell_index.hpp" #include "phlex/model/product_store.hpp" @@ -21,25 +20,13 @@ namespace { struct input_type_1 { int value; - auto operator<=>(input_type_1 const&) const = default; }; struct input_type_2 { int value; - auto operator<=>(input_type_2 const&) const = default; }; - struct output_type_1 { - int value; - auto operator<=>(output_type_1 const&) const = default; - }; - - output_type_1 add_values(input_type_1 const& left, input_type_2 const& right) - { - return {left.value + right.value}; - } - - template + template product_specification spec(char const* creator, char const* suffix) { return {algorithm_name{creator}, identifier{suffix}, make_type_id()}; @@ -58,57 +45,43 @@ namespace { store->add_product(spec(creator, suffix), std::move(value)); return store; } - - template - auto algorithm_bits_for(F function) - { - return algorithm_bits{std::shared_ptr{}, std::move(function)}; - } } -TEST_CASE("transform_node joins multiple input products", "[join]") +TEST_CASE("multilayer_join_node joins multiple input products", "[join]") { oneapi::tbb::flow::graph graph; - auto left_selector = selector("left_input", ""); - auto right_selector = selector("right_input", ""); auto left_store = store_with_product("left_input", "", input_type_1{17}); auto right_store = store_with_product("right_input", "", input_type_2{25}); - auto alg = algorithm_bits_for(add_values); - transform_node node{algorithm_name{"add_values"}, - 1u, - {}, - graph, - std::move(alg), - {left_selector, right_selector}, - {}}; - declared_transform& transform = node; + // Force repeaters by passing distinct layer names. + // The actual routing is performed by matching index hashes between data, index, and flush. + auto join = multilayer_join_node<2>{ + graph, + "multilayer_join_test", + std::vector{identifier{"left_layer"}, identifier{"right_layer"}}}; - oneapi::tbb::flow::queue_node sink{graph}; - make_edge(transform.output_port(), sink); + oneapi::tbb::flow::queue_node> sink{graph}; + make_edge(output_port<0>(join), sink); - auto& left_port = *transform.ports().at(0); - auto& right_port = *transform.ports().at(1); + auto& left_port = receiver_for<0ull, 2>(join, 0u); + auto& right_port = receiver_for<0ull, 2>(join, 1u); REQUIRE(left_port.try_put({.store = left_store, .id = message_id})); graph.wait_for_all(); - message output; + message_tuple<2> output; CHECK_FALSE(sink.try_get(output)); - CHECK(transform.num_calls() == 0u); REQUIRE(right_port.try_put({.store = right_store, .id = message_id})); graph.wait_for_all(); CHECK_FALSE(sink.try_get(output)); - CHECK(transform.num_calls() == 0u); - auto index_ports = transform.index_ports(); + auto index_ports = join.index_ports(); REQUIRE(index_ports.size() == 2u); REQUIRE(index_ports[0].index_port->try_put( {.index = left_store->index(), .msg_id = message_id, .cache = true})); graph.wait_for_all(); CHECK_FALSE(sink.try_get(output)); - CHECK(transform.num_calls() == 0u); REQUIRE(index_ports[1].index_port->try_put( {.index = right_store->index(), .msg_id = message_id, .cache = true})); @@ -117,19 +90,16 @@ TEST_CASE("transform_node joins multiple input products", "[join]") REQUIRE(sink.try_get(output)); CHECK_FALSE(sink.try_get(output)); + CHECK(std::get<0>(output).id == message_id); + CHECK(std::get<1>(output).id == message_id); + REQUIRE(std::get<0>(output).store); + REQUIRE(std::get<1>(output).store); + CHECK(std::get<0>(output).store->index() == left_store->index()); + CHECK(std::get<1>(output).store->index() == right_store->index()); + + // Do what is necessary to have the tokens flushed, so that the test does not generate + // warnings. REQUIRE(index_ports[0].token_port->try_put({.index = left_store->index(), .count = 1})); REQUIRE(index_ports[1].token_port->try_put({.index = right_store->index(), .count = 1})); graph.wait_for_all(); - - CHECK(output.id == message_id); - REQUIRE(output.store); - CHECK(output.store->index() == left_store->index()); - - auto const& output_specs = transform.output(); - REQUIRE(output_specs.size() == 1u); - CHECK(output_specs[0].suffix() == identifier{""}); - CHECK(output.store->get_product(output_specs[0]) == output_type_1{42}); - - CHECK(transform.num_calls() == 1u); - CHECK(transform.product_count() == 1u); } diff --git a/test/transform_node_test.cpp b/test/transform_node_test.cpp index 53c557c60..fff626ba6 100644 --- a/test/transform_node_test.cpp +++ b/test/transform_node_test.cpp @@ -41,11 +41,6 @@ namespace { output_type_1 double_value(input_type_1 const& input) { return {input.value * 2}; } - output_type_1 add_values(input_type_1 const& left, input_type_2 const& right) - { - return {left.value + right.value}; - } - auto number_and_label(input_type_1 const& input) { return std::tuple{output_type_1{input.value}, output_type_2{std::to_string(input.value)}}; @@ -114,75 +109,6 @@ TEST_CASE("transform_node directly transforms one input product", "[transform_no CHECK(transform.product_count() == 1u); } -TEST_CASE("transform_node joins multiple input products", "[transform_node]") -{ - oneapi::tbb::flow::graph graph; - auto left_selector = selector("left_input", ""); - auto right_selector = selector("right_input", ""); - auto left_store = store_with_product("left_input", "", input_type_1{17}); - auto right_store = store_with_product("right_input", "", input_type_2{25}); - auto alg = algorithm_bits_for(add_values); - - transform_node node{algorithm_name{"add_values"}, - 1u, - {}, - graph, - std::move(alg), - {left_selector, right_selector}, - {}}; - declared_transform& transform = node; - - auto const& output_specs = transform.output(); - REQUIRE(output_specs.size() == 1u); - CHECK(output_specs[0].suffix() == identifier{""}); - - oneapi::tbb::flow::queue_node sink{graph}; - make_edge(transform.output_port(), sink); - - auto& left_port = *transform.ports().at(0); - auto& right_port = *transform.ports().at(1); - - REQUIRE(left_port.try_put({.store = left_store, .id = message_id})); - graph.wait_for_all(); - - message output; - CHECK_FALSE(sink.try_get(output)); - CHECK(transform.num_calls() == 0u); - - REQUIRE(right_port.try_put({.store = right_store, .id = message_id})); - graph.wait_for_all(); - CHECK_FALSE(sink.try_get(output)); - CHECK(transform.num_calls() == 0u); - - auto index_ports = transform.index_ports(); - REQUIRE(index_ports.size() == 2u); - REQUIRE(index_ports[0].index_port->try_put( - {.index = left_store->index(), .msg_id = message_id, .cache = true})); - graph.wait_for_all(); - CHECK_FALSE(sink.try_get(output)); - CHECK(transform.num_calls() == 0u); - - REQUIRE(index_ports[1].index_port->try_put( - {.index = right_store->index(), .msg_id = message_id, .cache = true})); - graph.wait_for_all(); - - REQUIRE(sink.try_get(output)); - CHECK_FALSE(sink.try_get(output)); - - REQUIRE(index_ports[0].token_port->try_put({.index = left_store->index(), .count = 1})); - REQUIRE(index_ports[1].token_port->try_put({.index = right_store->index(), .count = 1})); - graph.wait_for_all(); - - CHECK(output.id == message_id); - REQUIRE(output.store); - CHECK(output.store->index() == left_store->index()); - - CHECK(output.store->get_product(output_specs[0]) == output_type_1{42}); - - CHECK(transform.num_calls() == 1u); - CHECK(transform.product_count() == 1u); -} - TEST_CASE("transform_node stores multiple output products", "[transform_node]") { oneapi::tbb::flow::graph graph; From ffcf08322be2e958d13668b8080d33d8a7cbe386 Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" <41898282+github-actions[bot]@users.noreply.github.com> Date: Fri, 31 Jul 2026 16:12:43 +0000 Subject: [PATCH 4/5] Apply clang-format fixes --- test/multilayer_join_test.cpp | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/test/multilayer_join_test.cpp b/test/multilayer_join_test.cpp index 1c04fd738..096d8f6e9 100644 --- a/test/multilayer_join_test.cpp +++ b/test/multilayer_join_test.cpp @@ -26,7 +26,7 @@ namespace { int value; }; - template + template product_specification spec(char const* creator, char const* suffix) { return {algorithm_name{creator}, identifier{suffix}, make_type_id()}; From 54a0ee5b7d975186bbe9d0bf0fd5e1e589a9f59b Mon Sep 17 00:00:00 2001 From: Marc Paterno Date: Fri, 31 Jul 2026 11:32:45 -0500 Subject: [PATCH 5/5] test(multilayer_join): remove unused selector helper --- test/multilayer_join_test.cpp | 6 ------ 1 file changed, 6 deletions(-) diff --git a/test/multilayer_join_test.cpp b/test/multilayer_join_test.cpp index 096d8f6e9..d8300f01e 100644 --- a/test/multilayer_join_test.cpp +++ b/test/multilayer_join_test.cpp @@ -32,12 +32,6 @@ namespace { return {algorithm_name{creator}, identifier{suffix}, make_type_id()}; } - template - product_selector selector(char const* creator, char const* suffix) - { - return {.creator = creator, .layer = "job", .suffix = suffix, .type = make_type_id()}; - } - template product_store_ptr store_with_product(char const* creator, char const* suffix, T value) {