#include #include #include #include #include #include #include #include #include #include #include #include #include namespace aethera::web::media { namespace { using Clock = std::chrono::steady_clock; void corun_until_or_throw(std::string_view operation, const std::function& completed) { std::string failure; std::exception_ptr runtime_failure; std::binary_semaphore finished{0}; auto coordinator = task_flow::make_task_graph("gallery.video.test.await"); const auto deadline = Clock::now() + std::chrono::seconds{10}; if (!coordinator->add("await", [&] { coordinator->corun_until([&] { if (completed()) return true; if (Clock::now() < deadline) return false; failure = std::string{operation} + " timed out"; return true; }); })) throw std::runtime_error("failed to create gallery test await task"); const auto submitted = coordinator->run([&](std::exception_ptr current_failure) noexcept { runtime_failure = current_failure; finished.release(); }); if (submitted != Run_Taskflow_Result::submitted) throw std::runtime_error("failed to submit gallery test await task"); /* Test-only join boundary: production paths remain completion-driven and never block a caller. */ finished.acquire(); if (runtime_failure) std::rethrow_exception(runtime_failure); if (!failure.empty()) throw std::runtime_error(failure); } void configure_frame_policy(const std::shared_ptr& control, std::string_view id) { ASSERT_TRUE(control->write_plot_property(id, "frame-analysis", "user_frames_per_second", 100.0).value("success", false)); ASSERT_NO_THROW(corun_until_or_throw("Plot policy", [&] { const auto diagnostics = control->plot_diagnostics(id); return diagnostics.value("available", false) && diagnostics.contains("frame_policy"); })); } void schedule_manual_frame(const std::shared_ptr& control, std::string_view id, std::uint64_t correlation_sequence) { const auto now = Clock::now(); control->request_frame(id, Gallery_Frame_Request{ .issued_at = now, .time_milliseconds = std::chrono::duration( now.time_since_epoch()).count(), .correlation_id = correlation_sequence, .width = 720, .height = 420 }); } TEST(Gallery_Video_Stream, Sink_Completion_Task_Releases_Each_Plots_Independent_Frame) { ASSERT_EQ(initialize_task_runtime(), Initialize_Task_Runtime_Result::initialized); const auto definitions = gallery_plot_definitions(); const auto found = std::ranges::find_if( definitions, [](const Plot_Definition& definition) { return definition.dimension == Plot_Dimension::two_d; }); ASSERT_NE(found, definitions.end()); const auto second = std::ranges::find_if( found + 1, definitions.end(), [](const Plot_Definition& definition) { return definition.dimension == Plot_Dimension::two_d; }); ASSERT_NE(second, definitions.end()); const std::string plot_a{found->id}; const std::string plot_b{second->id}; auto stream = Gallery_Video_Stream::create({{plot_a}, {plot_b}}); ASSERT_TRUE(stream); auto control = mcp::Control_Service::create(mcp::Gallery_Output{ .make_sink = [stream](std::string_view id) { return stream->frame_sink(id); }}); configure_frame_policy(control, plot_a); configure_frame_policy(control, plot_b); constexpr std::size_t frame_count{5}; std::atomic_size_t received{}; std::atomic> outstanding{}; const auto subscription = stream->subscribe( "gallery-video-test", [&](Gallery_Stream_Frame frame) { if (!frame.video) return false; received.fetch_add(1, std::memory_order_acq_rel); outstanding.store(std::move(frame.video), std::memory_order_release); return true; }, [] { return true; }, [] { return nlohmann::json::object(); }); const std::array schedule{ plot_a, plot_a, plot_b, plot_a, plot_b}; stream->request_taskflow_trace(1); for (std::size_t index = 0; index < schedule.size(); ++index) { const auto published_before = control->plot_diagnostics(schedule[index]).at("frame_policy").at("completed_frames").get(); const auto received_before = received.load(std::memory_order_acquire); schedule_manual_frame(control, schedule[index], index + 1U); ASSERT_NO_THROW(corun_until_or_throw("encoded Plot frame", [&] { return received.load(std::memory_order_acquire) > received_before; })); ASSERT_NO_THROW(corun_until_or_throw( "frame policy publication feedback", [&] { const auto diagnostics = control->plot_diagnostics(schedule[index]); return diagnostics.at("frame_policy").at("completed_frames") > published_before; })); outstanding.store({}, std::memory_order_release); } nlohmann::json diagnostics; ASSERT_NO_THROW(corun_until_or_throw("gallery media drain", [&] { auto current = stream->diagnostics(); if (current.at("pending_frame_count") != 0U || current.at("received_frame_count") != current.at("processed_frame_count") || current.at("received_frame_count") != current.at("encoded_frame_count")) return false; diagnostics = std::move(current); return true; })); EXPECT_GE(diagnostics.at("received_frame_count"), frame_count); EXPECT_EQ(diagnostics.at("processed_frame_count"), diagnostics.at("received_frame_count")); EXPECT_EQ(diagnostics.at("encoded_frame_count"), diagnostics.at("received_frame_count")); EXPECT_EQ(diagnostics.at("pending_frame_count"), 0U); ASSERT_NO_THROW(corun_until_or_throw("gallery Taskflow trace", [&] { return stream->taskflow_trace().value("complete", false); })); const auto trace = stream->taskflow_trace(); ASSERT_EQ(trace.at("frames").size(), 1U); const auto& graphs = trace.at("frames").front().at("graphs"); const auto graph = std::ranges::find_if(graphs, [](const auto& value) { return value.value("name", std::string{}) == "gallery.video.frame"; }); ASSERT_NE(graph, graphs.end()); const auto& nodes = graph->at("nodes"); const auto feedback = std::ranges::find_if(nodes, [](const auto& value) { return value.value("name", std::string{}) == "frame.policy.sink.complete"; }); ASSERT_NE(feedback, nodes.end()); EXPECT_EQ(feedback->at("attributes").at("owner"), "frame_policy"); stream->unsubscribe(subscription); stream->shutdown(); std::weak_ptr stream_lifetime = stream; stream.reset(); EXPECT_FALSE(stream_lifetime.expired()); std::atomic_bool stopped{}; control->stop([&stopped] { stopped.store(true, std::memory_order_release); }); ASSERT_NO_THROW(corun_until_or_throw("Frame Policy Sink retirement", [&] { return stopped.load(std::memory_order_acquire); })); EXPECT_TRUE(stream_lifetime.expired()); } } }