改方式

This commit is contained in:
2026-08-29 12:03:22 +08:00
parent 05a79fc351
commit 6045221174
18 changed files with 478 additions and 281 deletions
+168 -21
View File
@@ -1,4 +1,5 @@
#include "Gallery_Video_Stream.hpp"
#include <mcp/core/runtime/Taskflow_Trace_Json.hpp>
#include <frame_statistics.hpp>
#include <render_common.hpp>
#include "frame_sampling/Frame_Sampler.hpp"
@@ -57,13 +58,61 @@ void Gallery_Video_Stream::Private::initialize(
source_ids.push_back(plot.id);
sources.push_back(Source{std::move(plot)});
}
encode_source_slot = sources.size() - 1U;
sampler = std::make_unique<frame_sampling::Frame_Sampler>(
gallery_frame_rate, tile_width, tile_height, atlas_columns,
std::move(source_ids), encode_source_slot);
std::move(source_ids), sources.size() - 1U);
metric_completion_starts.resize(sources.size());
metric_rendered_starts.resize(sources.size());
}
bool Gallery_Video_Stream::Private::mark_taskflow_trace() {
auto remaining = taskflow_trace_remaining.load(std::memory_order_acquire);
while (remaining != 0) {
if (taskflow_trace_remaining.compare_exchange_weak(
remaining, remaining - 1, std::memory_order_acq_rel,
std::memory_order_acquire))
return true;
}
return false;
}
void Gallery_Video_Stream::Private::restore_taskflow_trace() noexcept {
taskflow_trace_remaining.fetch_add(1, std::memory_order_release);
}
void Gallery_Video_Stream::Private::store_taskflow_trace(
const Taskflow_Frame_Trace& value) {
const auto trace = std::make_shared<const nlohmann::json>(
taskflow_trace_json(value));
auto state = taskflow_trace_control.load(std::memory_order_acquire);
for (;;) {
const auto requested = static_cast<std::uint32_t>(state >> 32U);
const auto captured = static_cast<std::uint32_t>(state);
if (captured >= requested) return;
taskflow_trace_slots[captured].store(trace, std::memory_order_release);
const auto next = (static_cast<std::uint64_t>(requested) << 32U) |
static_cast<std::uint64_t>(captured + 1U);
if (taskflow_trace_control.compare_exchange_weak(
state, next, std::memory_order_release,
std::memory_order_acquire))
return;
}
}
nlohmann::json Gallery_Video_Stream::Private::taskflow_trace_response() const {
nlohmann::json frames = nlohmann::json::array();
const auto state = taskflow_trace_control.load(std::memory_order_acquire);
const auto requested = static_cast<std::uint32_t>(state >> 32U);
const auto captured = static_cast<std::uint32_t>(state);
for (std::uint32_t index = 0; index < captured; ++index)
if (const auto trace = taskflow_trace_slots[index].load(
std::memory_order_acquire))
frames.push_back(*trace);
const auto remaining = taskflow_trace_remaining.load(
std::memory_order_acquire);
return {
{"protocol", "aethera.taskflow.frames"}, {"version", 1},
{"requested", requested}, {"remaining", remaining},
{"captured", frames.size()},
{"complete", requested != 0 && frames.size() == requested},
{"frames", std::move(frames)}};
}
bool Gallery_Video_Stream::Private::consumer_accepts(
Gallery_Transport_Mode transport) const noexcept {
const auto current = consumers.load(std::memory_order_acquire);
@@ -86,11 +135,14 @@ bool Gallery_Video_Stream::Private::remove_consumer(
replacement->reserve(current->size() - 1);
std::ranges::copy_if(*current, std::back_inserter(*replacement),
[stream](const Consumer& value) { return value.id != stream; });
const bool became_empty = replacement->empty();
std::shared_ptr<const Consumers> desired{std::move(replacement)};
if (consumers.compare_exchange_weak(
current, std::move(desired), std::memory_order_acq_rel,
std::memory_order_acquire))
std::memory_order_acquire)) {
if (became_empty) sample_timer.cancel();
return true;
}
}
}
void Gallery_Video_Stream::Private::publish(
@@ -138,11 +190,7 @@ void Gallery_Video_Stream::Private::accept_frame(
if (stopping.load(std::memory_order_acquire) ||
failed.load(std::memory_order_acquire) || !frame)
return;
const auto result = sampler->accept_frame(slot, frame->pixels);
if (slot == encode_source_slot &&
result == frame_sampling::Frame_Sampler::Accept_Frame_Result::accepted) {
delivered_clock_ticks.fetch_add(1, std::memory_order_relaxed);
}
static_cast<void>(sampler->accept_frame(slot, frame->pixels));
}
void Gallery_Video_Stream::Private::update_metrics(
const Plot_Render_Tick& tick,
@@ -158,7 +206,7 @@ void Gallery_Video_Stream::Private::update_metrics(
metric_encoded_start = encoded_frame_count;
metric_pixel_start = pixel_frame_count;
metric_sample_start = sampler_state.sampled_frames;
metric_clock_start = delivered_clock_ticks.load(std::memory_order_relaxed);
metric_clock_start = sampling_clock_ticks.load(std::memory_order_relaxed);
for (std::size_t slot = 0; slot < composition.sources.size(); ++slot) {
metric_completion_starts[slot] =
composition.sources[slot].completion_count;
@@ -171,7 +219,7 @@ void Gallery_Video_Stream::Private::update_metrics(
if (elapsed < metric_interval) return;
const double seconds = std::chrono::duration<double>(elapsed).count();
const auto clock_total =
delivered_clock_ticks.load(std::memory_order_relaxed);
sampling_clock_ticks.load(std::memory_order_relaxed);
nlohmann::json source_metrics = nlohmann::json::object();
const auto description = sampler->describe();
for (std::size_t slot = 0; slot < composition.sources.size(); ++slot) {
@@ -199,10 +247,10 @@ void Gallery_Video_Stream::Private::update_metrics(
progress.rendered_correlation_id
},
{
"clock_lag_ticks",
progress.rendered_correlation_id != 0 &&
tick.sequence >= progress.rendered_correlation_id
? tick.sequence - progress.rendered_correlation_id
"frame_lag",
progress.rendered_sequence != 0 &&
progress.completion_sequence >= progress.rendered_sequence
? progress.completion_sequence - progress.rendered_sequence
: 0
}
};
@@ -395,12 +443,11 @@ void Gallery_Video_Stream::bind_plots() {
throw std::invalid_argument("gallery video stream requires a Plot");
/*
* 先把媒体尾部作为 encode source Plot 的 post-publish 外接 DAG 安装,
* 再 subscribe/start Plot。这样从第一帧开始 Taskflow trace 都能看到
* gallery.sample.capture 后分叉到 FFmpeg 与原始像素两个子模型,
* 不存在首帧已经构图后才追加 extension 的初始化窗口。
* 图集媒体有自己的固定采样时钟,不再依赖某一张 Plot 完成后才能运行。
* 每个来源只发布 latest;30 FPS tick 到达时若没有新完成帧或上一轮媒体
* DAG 尚未完成,直接合并,不排队也不占用 Worker 等待。
*/
auto media = std::make_unique<Task_Graph>("gallery.video.frame");
auto media = std::make_shared<Task_Graph>("gallery.video.frame");
auto sample = media->add("gallery.sample.capture", [weak] {
if (const auto owner = weak.lock()) {
auto& owner_data = static_cast<Private&>(*owner->d);
@@ -469,8 +516,76 @@ void Gallery_Video_Stream::bind_plots() {
publish_ffmpeg.precede(pack_pixels);
pack_pixels.precede(publish_pixels);
publish_pixels.precede(complete);
data.sources[data.encode_source_slot].entry.plot->attach_scene_completion(
std::move(media));
data.media_graph = media;
data.sample_timer = Frame_Scheduler::instance().make_timer(
[weak](Frame_Scheduler::Tick) {
const auto owner = weak.lock();
if (!owner) return;
auto& owner_data = static_cast<Private&>(*owner->d);
if (owner_data.stopping.load(std::memory_order_acquire) ||
owner_data.failed.load(std::memory_order_acquire))
return;
owner_data.sampling_clock_ticks.fetch_add(
1, std::memory_order_relaxed);
if (!owner_data.sampler->has_pending_frame()) return;
bool idle{};
if (!owner_data.media_busy.compare_exchange_strong(
idle, true, std::memory_order_acq_rel,
std::memory_order_acquire)) {
owner_data.skipped_sample_ticks.fetch_add(
1, std::memory_order_relaxed);
return;
}
const auto graph = owner_data.media_graph;
bool trace_reserved = owner_data.mark_taskflow_trace();
auto trace_frame = trace_reserved
? std::make_shared<Render_Frame>(Frame_Identity{
owner_data.sampling_clock_ticks.load(
std::memory_order_acquire), 0})
: std::shared_ptr<Render_Frame>{};
bool trace_started{};
if (trace_frame) {
trace_frame->request_taskflow_trace();
trace_started = aethera::detail::begin_taskflow_trace(
*trace_frame);
if (!trace_started) {
owner_data.restore_taskflow_trace();
trace_reserved = false;
trace_frame.reset();
}
}
try {
auto completion = [weak, graph, trace_frame, trace_started] {
static_cast<void>(graph);
if (trace_started)
aethera::detail::finish_taskflow_trace(*trace_frame);
if (const auto completed = weak.lock()) {
auto& completed_data =
static_cast<Private&>(*completed->d);
if (trace_started)
completed_data.store_taskflow_trace(
trace_frame->take_taskflow_trace());
completed_data.media_busy.store(
false, std::memory_order_release);
}
};
if (trace_frame)
aethera::detail::run_taskflow(
*graph, *trace_frame, "gallery.media",
std::move(completion));
else
aethera::detail::run_taskflow(
*graph, std::move(completion));
}
catch (...) {
if (trace_started)
aethera::detail::finish_taskflow_trace(*trace_frame);
if (trace_reserved) owner_data.restore_taskflow_trace();
owner_data.media_busy.store(false,
std::memory_order_release);
owner_data.fail(std::current_exception());
}
});
for (std::size_t slot = 0; slot < data.sources.size(); ++slot) {
auto& source = data.sources[slot];
@@ -512,7 +627,9 @@ Gallery_Video_Stream::Stream_Id Gallery_Video_Stream::subscribe(
std::make_shared<const Transport_Readiness>(std::move(readiness)),
std::make_shared<const Transport_Diagnostics>(std::move(diagnostics))};
auto current = data.consumers.load(std::memory_order_acquire);
bool activate_sampling{};
for (;;) {
activate_sampling = current->empty();
auto replacement = std::make_shared<Private::Consumers>();
replacement->reserve(current->size() + 1);
std::ranges::copy_if(*current, std::back_inserter(*replacement),
@@ -527,6 +644,8 @@ Gallery_Video_Stream::Stream_Id Gallery_Video_Stream::subscribe(
std::memory_order_acquire))
break;
}
if (activate_sampling)
data.sample_timer.start_periodic(gallery_frame_rate);
return id;
}
void Gallery_Video_Stream::unsubscribe(Stream_Id stream) {
@@ -542,6 +661,7 @@ void Gallery_Video_Stream::request_video_key_frame() {
void Gallery_Video_Stream::shutdown() noexcept {
auto& data = static_cast<Private&>(*d);
if (data.stopping.exchange(true, std::memory_order_acq_rel)) return;
data.sample_timer.cancel();
for (auto& source : data.sources) {
if (source.stream == 0) continue;
try {
@@ -609,4 +729,31 @@ nlohmann::json Gallery_Video_Stream::diagnostics() const {
}
return output;
}
void Gallery_Video_Stream::request_taskflow_trace(std::size_t frame_count) {
if (frame_count == 0 ||
frame_count > Private::maximum_taskflow_trace_frames)
throw std::invalid_argument(
"Gallery Taskflow trace frame_count must be between 1 and 120");
auto& data = static_cast<Private&>(*d);
auto control = data.taskflow_trace_control.load(std::memory_order_acquire);
for (;;) {
const auto requested = static_cast<std::uint32_t>(control >> 32U);
const auto captured = static_cast<std::uint32_t>(control);
if (requested != captured)
throw std::logic_error(
"A Gallery Taskflow trace request is already active");
const auto next = static_cast<std::uint64_t>(frame_count) << 32U;
if (data.taskflow_trace_control.compare_exchange_weak(
control, next, std::memory_order_release,
std::memory_order_acquire))
break;
}
for (auto& slot : data.taskflow_trace_slots)
slot.store({}, std::memory_order_release);
data.taskflow_trace_remaining.store(frame_count,
std::memory_order_release);
}
nlohmann::json Gallery_Video_Stream::taskflow_trace() const {
return static_cast<const Private&>(*d).taskflow_trace_response();
}
}
+2
View File
@@ -57,6 +57,8 @@ struct Gallery_Video_Stream : Def<Gallery_Video_Stream, Root>,
[[nodiscard]] std::string layout_description(
Gallery_Transport_Mode transport) const;
[[nodiscard]] nlohmann::json diagnostics() const;
void request_taskflow_trace(std::size_t frame_count);
[[nodiscard]] nlohmann::json taskflow_trace() const;
private:
void bind_plots();
};
+15 -2
View File
@@ -1,7 +1,9 @@
#pragma once
#include "frame_sampling/Frame_Sampler.hpp"
#include <Frame_Policy/Frame_Scheduler.hpp>
#include <frame_statistics.hpp>
#include <ownership.hpp>
#include <array>
#include <atomic>
#include <chrono>
namespace aethera::web {
@@ -28,6 +30,9 @@ struct Gallery_Video_Stream::Private : Prev_Private {
Sliding_Statistics publish_ms{600}; /* 两种传输发布统计计算器。 */
std::vector<Source> sources{}; /* 已按业务标识排序的稳定图集来源。 */
std::unique_ptr<frame_sampling::Frame_Sampler> sampler{}; /* latest 图像与 30 FPS deadline 的唯一状态源。 */
std::shared_ptr<Task_Graph> media_graph{}; /* 独立媒体时钟运行的持久 DAG;异步完成回调延长其生命周期。 */
Frame_Scheduler::Timer sample_timer{}; /* 与任意单张 Plot 完成速率解耦的 30 FPS 控制时钟。 */
std::atomic_bool media_busy{}; /* 单个媒体 DAG 的 latest-only 准入。 */
frame_sampling::ffmpeg::FFmpeg_Frame_Transport ffmpeg_transport; /* FFmpeg 子模型的唯一编码上下文。 */
std::optional<frame_sampling::Sampled_Gallery_Frame> active_sample{}; /* 当前 deadline 取得的原生图集视图。 */
std::shared_ptr<const Encoded_Video_Frame> active_video{}; /* 当前编码结果不可变共享所有权。 */
@@ -37,7 +42,7 @@ struct Gallery_Video_Stream::Private : Prev_Private {
std::make_shared<const Consumers>()}; /* 不可变订阅集合;写入使用 COW,发布只做一次原子读取。 */
std::atomic_bool stopping{}; /* 关闭开始后拒绝新工作。 */
std::atomic_uint64_t skipped_sample_ticks{}; /* 未到 deadline 或无可写消费者的采样触发数。 */
std::atomic_uint64_t delivered_clock_ticks{}; /* 已分发给 Plot 的时钟 tick。 */
std::atomic_uint64_t sampling_clock_ticks{}; /* 有消费者期间触发的独立媒体时钟 tick。 */
std::atomic_bool failed{}; /* 首次 Unknown Failure 后停止热路径。 */
std::optional<std::string> terminal_failure{}; /* 首次终止失败文本。 */
std::uint64_t encoded_frame_count{}; /* 累计编码成功帧数。 */
@@ -48,8 +53,12 @@ struct Gallery_Video_Stream::Private : Prev_Private {
std::uint64_t metric_clock_start{}; /* 窗口起点累计时钟数。 */
std::vector<std::uint64_t> metric_completion_starts{}; /* 各 Plot 窗口起点逻辑完成数。 */
std::vector<std::uint64_t> metric_rendered_starts{}; /* 各 Plot 窗口起点真实画面数。 */
std::size_t encode_source_slot{}; /* 承载媒体尾部 DAG 的 Plot 槽位。 */
std::uint64_t metric_pixel_start{}; /* 窗口起点累计像素帧数。 */
static constexpr std::size_t maximum_taskflow_trace_frames{120};
std::atomic_size_t taskflow_trace_remaining{}; /* 尚待捕获的真实媒体 DAG 次数。 */
std::atomic_uint64_t taskflow_trace_control{}; /* 高 32 位 requested,低 32 位 captured。 */
std::array<std::atomic<std::shared_ptr<const nlohmann::json>>,
maximum_taskflow_trace_frames> taskflow_trace_slots{};
Private();
~Private();
template <Attached Attached_Object>
@@ -79,5 +88,9 @@ struct Gallery_Video_Stream::Private : Prev_Private {
void publish_ffmpeg_frame();
void publish_pixel_frame();
void complete_media_frame();
[[nodiscard]] bool mark_taskflow_trace();
void restore_taskflow_trace() noexcept;
void store_taskflow_trace(const Taskflow_Frame_Trace& trace);
[[nodiscard]] nlohmann::json taskflow_trace_response() const;
};
}
+11
View File
@@ -1,6 +1,7 @@
#include "Gallery_WebSocket.hpp"
#include "WebRtc_Video_Session.hpp"
#include <nlohmann/json.hpp>
#include <trantor/net/EventLoop.h>
#include <atomic>
#include <chrono>
#include <exception>
@@ -61,7 +62,17 @@ void Gallery_WebSocket::start() {
if (d->attached.exchange(true, std::memory_order_acq_rel)) return;
const auto weak = weak_from_this();
if (d->transport == Gallery_Transport_Mode::ffmpeg) {
const auto connection = d->connection.lock();
if (!connection)
throw std::logic_error(
"WebRTC session requires a live WebSocket connection");
auto* const owner_loop =
trantor::EventLoop::getEventLoopOfCurrentThread();
if (!owner_loop)
throw std::logic_error(
"WebRTC session must start on a Drogon I/O thread");
d->video = std::make_unique<WebRtc_Video_Session>(
*owner_loop,
[weak](std::string signal) {
const auto socket = weak.lock();
if (!socket) return;
+73 -45
View File
@@ -1,7 +1,9 @@
#include "WebRtc_Video_Session.hpp"
#include <ownership.hpp>
#include <concurrentqueue-1.0.5/blockingconcurrentqueue.h>
#include <nlohmann/json.hpp>
#include <rtc/rtc.hpp>
#include <trantor/net/EventLoop.h>
#include <atomic>
#include <chrono>
#include <cstddef>
@@ -35,7 +37,8 @@ void update_maximum(std::atomic_uint64_t& maximum,
}
}
struct WebRtc_Video_Session::Private {
struct WebRtc_Video_Session::Private final
: std::enable_shared_from_this<WebRtc_Video_Session::Private> {
struct Callback_State {
Signal_Handler signal_handler; /* SDP/ICE 信令的唯一交付出口。 */
Ready_Handler ready_handler; /* Track 打开后请求首个 IDR。 */
@@ -78,13 +81,14 @@ struct WebRtc_Video_Session::Private {
std::atomic_size_t transport_buffered_bytes{};
std::atomic_uint64_t readiness_check_count{};
std::atomic_uint64_t readiness_reject_count{};
std::atomic_uint64_t close_join_started_ns{};
std::atomic_uint64_t close_join_total_ns{};
std::atomic_uint64_t close_join_max_ns{};
std::atomic_bool close_join_active{};
Private(Signal_Handler signal_handler, Ready_Handler ready_handler,
not_null<trantor::EventLoop*> owner_loop; /* Peer/Track 创建所在 Drogon I/O 线程;服务停止前始终有效。 */
std::atomic_bool retirement_requested{}; /* close 已撤销入口并要求发送线程退出。 */
std::atomic_bool retirement_dispatched{}; /* 所属 I/O 线程只接收一次最终资源退役。 */
Private(trantor::EventLoop& value_owner_loop,
Signal_Handler signal_handler, Ready_Handler ready_handler,
Failure_Handler failure_handler)
: callbacks(std::make_shared<Callback_State>()) {
: callbacks(std::make_shared<Callback_State>()),
owner_loop(&value_owner_loop) {
if (!signal_handler || !failure_handler)
throw std::invalid_argument(
"WebRTC session requires signal and failure handlers");
@@ -93,12 +97,54 @@ struct WebRtc_Video_Session::Private {
callbacks->failure_handler = std::move(failure_handler);
}
void release_transport() noexcept {
try {
std::lock_guard lock(media_mutex);
const auto track = video_track.exchange({}, std::memory_order_acq_rel);
try { if (track) track->close(); }
catch (...) {}
try { if (peer) peer->close(); }
catch (...) {}
rtp_config.reset();
peer.reset();
}
catch (...) {}
}
void dispatch_retirement() noexcept {
if (!retirement_requested.load(std::memory_order_acquire) ||
sender_loop_active.load(std::memory_order_acquire))
return;
bool expected{};
if (!retirement_dispatched.compare_exchange_strong(
expected, true, std::memory_order_acq_rel,
std::memory_order_acquire))
return;
/* sender_thread 已退出或正处于返回尾声;detach 只解除 std::thread 句柄,
* 不等待。shared ownership 保证 Private 活到所属 I/O 线程完成销毁。 */
if (sender_thread.joinable()) sender_thread.detach();
try {
owner_loop->queueInLoop([self = shared_from_this()] {
self->release_transport();
});
}
catch (...) {
/* 无法把线程亲和资源交回所属循环属于不可恢复的生命周期故障;
* 禁止在错误线程直接销毁,也不静默回退为同步 join。 */
std::terminate();
}
}
void run_sender() noexcept {
sender_loop_active.store(true, std::memory_order_release);
struct Sender_Exit final {
std::atomic_bool& active;
~Sender_Exit() { active.store(false, std::memory_order_release); }
} sender_exit{sender_loop_active};
not_null<Private*> owner;
~Sender_Exit() {
owner->sender_loop_active.store(false, std::memory_order_release);
owner->dispatch_retirement();
}
} sender_exit{this};
for (;;) {
Encoded_Video_Frame frame;
pending_frames.wait_dequeue(frame);
@@ -151,10 +197,11 @@ struct WebRtc_Video_Session::Private {
};
WebRtc_Video_Session::WebRtc_Video_Session(
Signal_Handler signal_handler, Ready_Handler ready_handler,
trantor::EventLoop& owner_loop, Signal_Handler signal_handler,
Ready_Handler ready_handler,
Failure_Handler failure_handler)
: d(std::make_unique<Private>(
std::move(signal_handler), std::move(ready_handler),
: d(std::make_shared<Private>(
owner_loop, std::move(signal_handler), std::move(ready_handler),
std::move(failure_handler))) {}
WebRtc_Video_Session::~WebRtc_Video_Session() { close(); }
@@ -254,7 +301,14 @@ void WebRtc_Video_Session::start() {
d->peer = std::move(peer);
d->video_track.store(std::move(video_track), std::memory_order_release);
d->rtp_config = std::move(rtp_config);
d->sender_thread = std::thread([data = d.get()] { data->run_sender(); });
d->sender_loop_active.store(true, std::memory_order_release);
try {
d->sender_thread = std::thread([data = d] { data->run_sender(); });
}
catch (...) {
d->sender_loop_active.store(false, std::memory_order_release);
throw;
}
}
void WebRtc_Video_Session::accept_answer(std::string_view sdp) {
@@ -342,8 +396,6 @@ nlohmann::json WebRtc_Video_Session::diagnostics() const {
const auto send_total_ns = d->send_total_ns.load(std::memory_order_relaxed);
const auto readiness_check_count =
d->readiness_check_count.load(std::memory_order_relaxed);
const auto close_join_total_ns =
d->close_join_total_ns.load(std::memory_order_relaxed);
const auto milliseconds = [](std::uint64_t nanoseconds) {
return static_cast<double>(nanoseconds) / 1'000'000.0;
};
@@ -370,11 +422,7 @@ nlohmann::json WebRtc_Video_Session::diagnostics() const {
{"current_send_ms", milliseconds(d->send_active.load(std::memory_order_acquire)
? now - d->send_started_ns.load(std::memory_order_relaxed) : 0)},
{"readiness_check_count", readiness_check_count},
{"readiness_reject_count", d->readiness_reject_count.load(std::memory_order_relaxed)},
{"close_join_total_ms", milliseconds(close_join_total_ns)},
{"close_join_max_ms", milliseconds(d->close_join_max_ns.load(std::memory_order_relaxed))},
{"current_close_join_ms", milliseconds(d->close_join_active.load(std::memory_order_acquire)
? now - d->close_join_started_ns.load(std::memory_order_relaxed) : 0)}};
{"readiness_reject_count", d->readiness_reject_count.load(std::memory_order_relaxed)}};
}
void WebRtc_Video_Session::close() noexcept {
@@ -382,32 +430,12 @@ void WebRtc_Video_Session::close() noexcept {
if (d->callbacks->closed.exchange(true, std::memory_order_acq_rel))
return;
d->callbacks->track_open.store(false, std::memory_order_release);
d->retirement_requested.store(true, std::memory_order_release);
Encoded_Video_Frame discarded;
while (d->pending_frames.try_dequeue(discarded)) {}
while (!d->pending_frames.try_enqueue(
Encoded_Video_Frame{}))
std::this_thread::yield();
if (d->sender_thread.joinable()) {
const auto join_started = steady_nanoseconds();
d->close_join_started_ns.store(join_started, std::memory_order_relaxed);
d->close_join_active.store(true, std::memory_order_release);
d->sender_thread.join();
const auto join_elapsed = steady_nanoseconds() - join_started;
d->close_join_active.store(false, std::memory_order_release);
d->close_join_total_ns.fetch_add(join_elapsed, std::memory_order_relaxed);
update_maximum(d->close_join_max_ns, join_elapsed);
}
while (d->pending_frames.try_dequeue(discarded)) {}
d->outstanding_video_frames.store(0, std::memory_order_release);
std::lock_guard lock(d->media_mutex);
const auto track = d->video_track.exchange({}, std::memory_order_acq_rel);
try { if (track) track->close(); }
catch (...) {}
try { if (d->peer) d->peer->close(); }
catch (...) {}
d->rtp_config.reset();
d->peer.reset();
if (!d->pending_frames.try_enqueue(Encoded_Video_Frame{}))
std::terminate();
d->dispatch_retirement();
}
catch (...) {}
}
+6 -2
View File
@@ -7,6 +7,9 @@
#include <nlohmann/json_fwd.hpp>
#include <string>
#include <string_view>
namespace trantor {
class EventLoop;
}
namespace aethera::web {
struct WebRtc_Video_Session final {
public:
@@ -17,7 +20,8 @@ public:
using Signal_Handler = std::function<void(std::string)>;
using Ready_Handler = std::function<void()>;
using Failure_Handler = std::function<void(std::exception_ptr)>;
WebRtc_Video_Session(Signal_Handler signal_handler,
WebRtc_Video_Session(trantor::EventLoop& owner_loop,
Signal_Handler signal_handler,
Ready_Handler ready_handler,
Failure_Handler failure_handler);
~WebRtc_Video_Session();
@@ -33,6 +37,6 @@ public:
void close() noexcept;
private:
struct Private;
std::unique_ptr<Private> d;
std::shared_ptr<Private> d;
};
}
+20 -21
View File
@@ -53,10 +53,8 @@ bool taskflow_graph_contains_gallery_media(const nlohmann::json& graph) {
}
/*
* Gallery 的媒体 DAG 实际只在组内唯一 encode-source Plot 上执行。诊断接口
* 在服务端把该真实执行帧中由 Task_Graph 包装器捕获到的媒体 graph 附加到
* 当前 Plot 的 trace 响应;不复制执行、不伪造 FFmpeg 节点,也不让前端知道
* encode-source 是哪个 Plot。
* Gallery 的媒体 DAG 由媒体组自己的固定时钟执行。诊断接口只把该媒体组
* 实际捕获的 graph 附加到当前 Plot trace;不复制执行,也不伪造节点。
*/
nlohmann::json merge_gallery_media_trace(nlohmann::json plot_trace,
const nlohmann::json& media_trace) {
@@ -134,17 +132,16 @@ int run_web_server(std::uint16_t port, const std::filesystem::path& asset_root)
auto gallery_streams = std::make_shared<Gallery_Stream_Map>();
auto plot_media = std::make_shared<
std::unordered_map<std::string, std::string>>();
auto plot_media_source = std::make_shared<
auto plot_media_group = std::make_shared<
std::unordered_map<std::string, std::string>>();
const auto add_media_group = [&gallery_streams, &plot_media, &plot_media_source](
const auto add_media_group = [&gallery_streams, &plot_media, &plot_media_group](
std::string id,
std::vector<Gallery_Video_Stream::Plot_Entry> entries) {
if (entries.empty()) return;
const auto path = "/ws/gallery?group=" + id;
const auto media_source = entries.back().id;
for (const auto& entry : entries) {
plot_media->emplace(entry.id, path);
plot_media_source->emplace(entry.id, media_source);
plot_media_group->emplace(entry.id, id);
}
gallery_streams->emplace(
std::move(id), Gallery_Video_Stream::create(std::move(entries)));
@@ -237,7 +234,8 @@ int run_web_server(std::uint16_t port, const std::filesystem::path& asset_root)
callback(json_response(output.content));
}, {drogon::Get});
app.registerHandler("/plot/{1}/taskflow", [control, plot_media_source](
app.registerHandler("/plot/{1}/taskflow", [control, plot_media_group,
gallery_streams](
const drogon::HttpRequestPtr& request,
std::function<void(const drogon::HttpResponsePtr&)>&& callback,
std::string plot_id) {
@@ -246,10 +244,16 @@ int run_web_server(std::uint16_t port, const std::filesystem::path& asset_root)
callback(error_response(drogon::k404NotFound, "unknown plot"));
return;
}
const auto source_found = plot_media_source->find(plot_id);
auto media_plot = source_found == plot_media_source->end()
? plot : control->find_plot(source_found->second);
if (!media_plot) media_plot = plot;
const auto group_found = plot_media_group->find(plot_id);
const auto stream_found = group_found == plot_media_group->end()
? gallery_streams->end()
: gallery_streams->find(group_found->second);
if (stream_found == gallery_streams->end()) {
callback(error_response(drogon::k404NotFound,
"plot has no gallery media group"));
return;
}
const auto& media_stream = stream_found->second;
try {
if (request->method() == drogon::Post) {
const auto input = nlohmann::json::parse(request->body());
@@ -265,19 +269,14 @@ int run_web_server(std::uint16_t port, const std::filesystem::path& asset_root)
state.value("captured", std::size_t{});
};
if (trace_active(plot->taskflow_trace()) ||
trace_active(media_plot->post_publish_taskflow_trace()))
trace_active(media_stream->taskflow_trace()))
throw std::logic_error(
"A Taskflow frame trace request is already active");
/*
* 非 encode-source Plot 只让真实媒体源捕获 post-publish DAG。
* 不再额外追踪另一张 Plot 的 Render/Prepare/Paint,避免诊断本身
* 放大 Executor 压力。FFmpeg 仍来自真实 Task_Graph 包装节点。
*/
media_plot->request_post_publish_taskflow_trace(frame_count);
media_stream->request_taskflow_trace(frame_count);
plot->request_taskflow_trace(frame_count);
}
auto output = plot->taskflow_trace();
const auto media = media_plot->post_publish_taskflow_trace();
const auto media = media_stream->taskflow_trace();
output = merge_gallery_media_trace(std::move(output), media);
const bool plot_complete = output.value("complete", false);
const bool media_complete = media.value("complete", false);
@@ -22,6 +22,7 @@ struct Frame_Sampler::Private {
std::size_t sampling_source_slot{}; /* 为样本提供页面媒体时间的固定业务来源。 */
std::atomic_int64_t next_deadline_ns{}; /* 下一次允许生成样本的单调 deadline。 */
std::atomic_uint64_t accepted_source_frames{}; /* 成功发布到 latest 的源帧计数。 */
std::atomic_uint64_t sampled_source_generation{}; /* 最近样本覆盖到的 accepted_source_frames。 */
std::atomic_uint64_t rejected_source_frames{}; /* 无效或过期源帧计数。 */
std::atomic_uint64_t sampled_frames{}; /* 实际采样计数。 */
std::atomic_uint64_t early_ticks{}; /* deadline 前被合并的触发计数。 */
@@ -63,7 +64,7 @@ Frame_Sampler::Accept_Frame_Result Frame_Sampler::accept_frame(
std::size_t slot, std::shared_ptr<const Plot_Pixel_Frame> frame) {
const auto result = d->atlas.accept_frame(slot, std::move(frame));
if (result == detail::Gallery_Frame_Atlas::Accept_Frame_Result::accepted) {
d->accepted_source_frames.fetch_add(1, std::memory_order_relaxed);
d->accepted_source_frames.fetch_add(1, std::memory_order_release);
return Accept_Frame_Result::accepted;
}
d->rejected_source_frames.fetch_add(1, std::memory_order_relaxed);
@@ -72,10 +73,16 @@ Frame_Sampler::Accept_Frame_Result Frame_Sampler::accept_frame(
: Accept_Frame_Result::invalid_frame;
}
bool Frame_Sampler::has_pending_frame() const noexcept {
return d->accepted_source_frames.load(std::memory_order_acquire) !=
d->sampled_source_generation.load(std::memory_order_acquire);
}
std::optional<Sampled_Gallery_Frame> Frame_Sampler::sample() {
const auto now = Clock::now();
const auto now_ns = clock_nanoseconds(now);
const auto compose_sample = [this, now](double deadline_delay_ms) {
const auto compose_sample = [this, now](double deadline_delay_ms,
std::uint64_t generation) {
auto composition = d->atlas.compose();
const auto& source = composition.sources[d->sampling_source_slot];
Plot_Render_Tick tick{
@@ -85,9 +92,16 @@ std::optional<Sampled_Gallery_Frame> Frame_Sampler::sample() {
1'000.0,
composition.width,
composition.height};
d->sampled_source_generation.store(generation,
std::memory_order_release);
return Sampled_Gallery_Frame{
std::move(composition), tick, deadline_delay_ms};
};
const auto generation = d->accepted_source_frames.load(
std::memory_order_acquire);
if (generation == d->sampled_source_generation.load(
std::memory_order_acquire))
return std::nullopt;
auto deadline_ns = d->next_deadline_ns.load(std::memory_order_acquire);
for (;;) {
if (deadline_ns == 0) {
@@ -97,7 +111,7 @@ std::optional<Sampled_Gallery_Frame> Frame_Sampler::sample() {
std::memory_order_acquire))
continue;
d->sampled_frames.fetch_add(1, std::memory_order_relaxed);
return compose_sample(0.0);
return compose_sample(0.0, generation);
}
if (now_ns < deadline_ns) {
d->early_ticks.fetch_add(1, std::memory_order_relaxed);
@@ -114,7 +128,8 @@ std::optional<Sampled_Gallery_Frame> Frame_Sampler::sample() {
continue;
d->sampled_frames.fetch_add(1, std::memory_order_relaxed);
d->missed_periods.fetch_add(missed, std::memory_order_relaxed);
return compose_sample(static_cast<double>(late_ns) / 1'000'000.0);
return compose_sample(static_cast<double>(late_ns) / 1'000'000.0,
generation);
}
}
@@ -50,6 +50,7 @@ public:
[[nodiscard]] Accept_Frame_Result accept_frame(
std::size_t slot, std::shared_ptr<const Plot_Pixel_Frame> frame);
[[nodiscard]] bool has_pending_frame() const noexcept;
[[nodiscard]] std::optional<Sampled_Gallery_Frame> sample();
[[nodiscard]] detail::Gallery_Atlas_Description describe() const;
[[nodiscard]] Frame_Sampler_State state() const noexcept;
+3 -1
View File
@@ -24,8 +24,10 @@ TEST(Frame_Sampler, First_Deadline_Samples_Latest_Without_Format_Conversion) {
Frame_Sampler sampler(30.0, 2, 2, 1, {"plot"}, 0);
ASSERT_EQ(sampler.accept_frame(0, native_bgra_frame()),
Frame_Sampler::Accept_Frame_Result::accepted);
EXPECT_TRUE(sampler.has_pending_frame());
const auto sample = sampler.sample();
ASSERT_TRUE(sample);
EXPECT_FALSE(sampler.has_pending_frame());
EXPECT_EQ(sample->composition.layout, Plot_Pixel_Layout::bgra8);
EXPECT_EQ(sample->composition.pixels[0], std::byte{11});
EXPECT_EQ(sample->composition.pixels[1], std::byte{22});
@@ -37,7 +39,7 @@ TEST(Frame_Sampler, First_Deadline_Samples_Latest_Without_Format_Conversion) {
const auto state = sampler.state();
EXPECT_DOUBLE_EQ(state.target_frame_rate_fps, 30.0);
EXPECT_EQ(state.sampled_frames, 1U);
EXPECT_EQ(state.early_ticks, 1U);
EXPECT_EQ(state.early_ticks, 0U);
}
TEST(WebSocket_Pixel_Frame, Header_Describes_Unchanged_Native_Payload) {