优化了 还是卡

This commit is contained in:
2026-08-25 21:44:57 +08:00
parent 4f3430e3ad
commit 65799cfc2a
21 changed files with 850 additions and 1224 deletions
+46 -58
View File
@@ -10,8 +10,6 @@
#include <chrono>
#include <cmath>
#include <future>
#include <mutex>
#include <shared_mutex>
#include <optional>
#include <stdexcept>
#include <string>
@@ -123,8 +121,7 @@ struct Gallery_Video_Stream::Private {
Sliding_Statistics compose_ms{600};
Sliding_Statistics encode_ms{600};
Sliding_Statistics publish_ms{600};
mutable std::shared_mutex diagnostics_exchange_mutex;
double_buffer::Double_Buffer<nlohmann::json> diagnostics{};
std::atomic<std::shared_ptr<const nlohmann::json>> diagnostics{};
} state{};
std::vector<Source> sources{}; /* 已按业务标识排序的稳定图集来源。 */
@@ -141,9 +138,8 @@ struct Gallery_Video_Stream::Private {
std::shared_ptr<const Video_Readiness> video_readiness{}; /* 编码前查询 WebRTC 三槽发送窗口。 */
std::shared_ptr<const Video_Diagnostics> video_diagnostics{};
};
mutable std::mutex consumers_mutex;
Stream_Id next_consumer_id{1}; /* 在互斥量内区分重连前后的独占消费者。 */
std::optional<Consumer> consumer{}; /* 独占网页对应的唯一媒体消费者。 */
std::atomic_uint64_t next_consumer_id{1};
std::atomic<std::shared_ptr<const Consumer>> consumer{}; /* 最新连接原子接管旧连接。 */
std::atomic_bool stopping{}; /* 拒绝关闭开始后的新 tick、帧与订阅。 */
std::atomic_bool key_frame_requested{}; /* 下一次实际编码前消费的关键帧请求。 */
std::atomic_uint64_t skipped_encode_ticks{}; /* 被最新帧策略合并的页面编码 tick 总数。 */
@@ -177,16 +173,13 @@ struct Gallery_Video_Stream::Private {
[[nodiscard]] std::shared_ptr<const Stream_Handler>
consumer_handler() const {
std::lock_guard lock(consumers_mutex);
return consumer ? consumer->handler : nullptr;
const auto current = consumer.load(std::memory_order_acquire);
return current ? current->handler : nullptr;
}
[[nodiscard]] bool consumer_accepts_video() const noexcept {
std::shared_ptr<const Video_Readiness> readiness;
{
std::lock_guard lock(consumers_mutex);
if (consumer) readiness = consumer->video_readiness;
}
const auto current = consumer.load(std::memory_order_acquire);
const auto readiness = current ? current->video_readiness : nullptr;
if (!readiness) return false;
try { return (*readiness)(); }
catch (...) { return false; }
@@ -204,11 +197,11 @@ struct Gallery_Video_Stream::Private {
return;
}
catch (...) {}
{
std::lock_guard lock(consumers_mutex);
if (consumer && consumer->handler == handler)
consumer.reset();
}
auto current = consumer.load(std::memory_order_acquire);
if (current && current->handler == handler)
static_cast<void>(consumer.compare_exchange_strong(
current, {}, std::memory_order_acq_rel,
std::memory_order_acquire));
if (clock) clock->stop();
}
catch (...) {}
@@ -345,9 +338,9 @@ struct Gallery_Video_Stream::Private {
metric_started = now;
metric_encoded_start = encoded_frame_count;
metric_clock_start = clock_total;
*state.diagnostics.current = std::move(output);
std::unique_lock lock(state.diagnostics_exchange_mutex);
state.diagnostics.advance();
state.diagnostics.store(
std::make_shared<const nlohmann::json>(std::move(output)),
std::memory_order_release);
}
void compose_media_frame() {
@@ -502,25 +495,27 @@ Gallery_Video_Stream::Stream_Id Gallery_Video_Stream::subscribe(
if (!video_diagnostics)
throw std::invalid_argument(
"gallery video subscription requires WebRTC diagnostics");
Stream_Id id{};
{
std::lock_guard lock(d->consumers_mutex);
id = d->next_consumer_id++;
// 每个媒体组只保留最新浏览器连接;新页面原子替换旧回调,旧连接随后
// unsubscribe(id) 时不会误删新连接。连接接管不依赖页面租约或关闭时序。
d->consumer = Private::Consumer{
const auto id = d->next_consumer_id.fetch_add(
1, std::memory_order_relaxed);
// 每个媒体组只保留最新浏览器连接;新页面原子替换旧回调,旧连接随后
// unsubscribe(id) 时不会误删新连接。连接接管不依赖页面租约或关闭时序。
auto consumer = std::make_shared<const Private::Consumer>(
Private::Consumer{
id, std::make_shared<const Stream_Handler>(std::move(handler)),
std::make_shared<const Video_Readiness>(
std::move(video_readiness)),
std::make_shared<const Video_Diagnostics>(
std::move(video_diagnostics))};
}
std::move(video_diagnostics))});
d->consumer.store(consumer, std::memory_order_release);
try {
d->clock->start();
}
catch (...) {
std::lock_guard lock(d->consumers_mutex);
if (d->consumer && d->consumer->id == id) d->consumer.reset();
auto current = d->consumer.load(std::memory_order_acquire);
if (current && current->id == id)
static_cast<void>(d->consumer.compare_exchange_strong(
current, {}, std::memory_order_acq_rel,
std::memory_order_acquire));
throw;
}
return id;
@@ -528,11 +523,13 @@ Gallery_Video_Stream::Stream_Id Gallery_Video_Stream::subscribe(
void Gallery_Video_Stream::unsubscribe(Stream_Id stream) {
bool removed{};
{
std::lock_guard lock(d->consumers_mutex);
if (d->consumer && d->consumer->id == stream) {
d->consumer.reset();
auto current = d->consumer.load(std::memory_order_acquire);
while (current && current->id == stream) {
if (d->consumer.compare_exchange_weak(
current, {}, std::memory_order_acq_rel,
std::memory_order_acquire)) {
removed = true;
break;
}
}
if (removed) d->clock->stop();
@@ -553,11 +550,7 @@ void Gallery_Video_Stream::shutdown() noexcept {
catch (...) {}
source.stream = 0;
}
try {
std::lock_guard lock(d->consumers_mutex);
d->consumer.reset();
}
catch (...) {}
d->consumer.store({}, std::memory_order_release);
}
std::string Gallery_Video_Stream::layout_description() const {
@@ -583,22 +576,17 @@ std::string Gallery_Video_Stream::layout_description() const {
}
nlohmann::json Gallery_Video_Stream::diagnostics() const {
nlohmann::json output;
{
std::shared_lock lock(d->state.diagnostics_exchange_mutex);
output = !d->state.diagnostics.pending->is_null()
? *d->state.diagnostics.pending
: nlohmann::json{{"kind", "gallery_metrics"},
{"protocol", "aethera.gallery.video"},
{"version", 2},
{"sources", nlohmann::json::object()}};
}
std::shared_ptr<const Video_Diagnostics> video_diagnostics;
{
std::lock_guard lock(d->consumers_mutex);
if (d->consumer)
video_diagnostics = d->consumer->video_diagnostics;
}
const auto diagnostics = d->state.diagnostics.load(
std::memory_order_acquire);
nlohmann::json output = diagnostics
? *diagnostics
: nlohmann::json{{"kind", "gallery_metrics"},
{"protocol", "aethera.gallery.video"},
{"version", 2},
{"sources", nlohmann::json::object()}};
const auto current = d->consumer.load(std::memory_order_acquire);
const auto video_diagnostics = current
? current->video_diagnostics : nullptr;
if (video_diagnostics) output["webrtc"] = (*video_diagnostics)();
return output;
}
+6 -53
View File
@@ -27,8 +27,6 @@ std::string media_failure_description(const std::exception_ptr& failure) {
struct Gallery_WebSocket::Private {
std::weak_ptr<drogon::WebSocketConnection> connection; /* 仅在连接存活时投递信令。 */
std::shared_ptr<Gallery_Video_Stream> stream; /* 当前媒体组共享的编码图集。 */
std::shared_ptr<Page_Session> page_session; /* 当前连接所属的唯一页面权威。 */
Page_Session::Connection_Id page_connection{}; /* 在页面权威内登记的连接标识。 */
std::unique_ptr<WebRtc_Video_Session> video; /* 当前浏览器连接的 WebRTC 发送会话。 */
Gallery_Video_Stream::Stream_Id subscription{}; /* 共享图集输出的订阅标识。 */
std::atomic_bool attached{}; /* start 成功后为真,并保证 close 只执行一次。 */
@@ -37,14 +35,10 @@ struct Gallery_WebSocket::Private {
Gallery_WebSocket::Gallery_WebSocket(
drogon::WebSocketConnectionPtr connection,
std::shared_ptr<Gallery_Video_Stream> stream,
std::shared_ptr<Page_Session> page_session,
Page_Session::Connection_Id page_connection)
std::shared_ptr<Gallery_Video_Stream> stream)
: d(std::make_unique<Private>()) {
d->connection = std::move(connection);
d->stream = std::move(stream);
d->page_session = std::move(page_session);
d->page_connection = page_connection;
}
Gallery_WebSocket::~Gallery_WebSocket() { close(); }
@@ -150,8 +144,6 @@ void Gallery_WebSocket::receive(std::string_view message) {
void Gallery_WebSocket::close() noexcept {
if (!d->attached.exchange(false, std::memory_order_acq_rel)) return;
try { d->page_session->detach(d->page_connection); }
catch (...) {}
try {
if (d->subscription != 0) d->stream->unsubscribe(d->subscription);
}
@@ -161,13 +153,11 @@ void Gallery_WebSocket::close() noexcept {
}
Gallery_WebSocket_Controller::Gallery_WebSocket_Controller(
Stream_Resolver value_resolver,
std::shared_ptr<Page_Session> value_page_session)
: resolve_stream(std::move(value_resolver)),
page_session(std::move(value_page_session)) {
if (!resolve_stream || !page_session)
Stream_Resolver value_resolver)
: resolve_stream(std::move(value_resolver)) {
if (!resolve_stream)
throw std::invalid_argument(
"gallery WebSocket requires media and page-session services");
"gallery WebSocket requires a media stream resolver");
}
void Gallery_WebSocket_Controller::initPathRouting() {
@@ -179,43 +169,8 @@ void Gallery_WebSocket_Controller::handleNewConnection(
const drogon::HttpRequestPtr& request,
const drogon::WebSocketConnectionPtr& connection) {
try {
const auto page_connection = page_session->attach(
request->getParameter("session"), [weak = std::weak_ptr(connection)] {
const auto active = weak.lock();
if (!active || !active->connected()) return;
try {
active->send(nlohmann::json{
{"kind", "page_session_replaced"},
{"protocol", "aethera.page.session"},
{"version", 1}}.dump(),
drogon::WebSocketMessageType::Text);
}
catch (...) {}
try {
active->shutdown(drogon::CloseCode::kViolation,
"Page session replaced");
}
catch (...) {}
});
if (!page_connection) {
try {
connection->send(nlohmann::json{
{"kind", "page_session_rejected"},
{"protocol", "aethera.page.session"},
{"version", 1}}.dump(),
drogon::WebSocketMessageType::Text);
}
catch (...) {}
try {
connection->shutdown(drogon::CloseCode::kViolation,
"Stale page session");
}
catch (...) {}
return;
}
const auto stream = resolve_stream(request->getParameter("group"));
if (!stream) {
page_session->detach(*page_connection);
try {
connection->send(nlohmann::json{
{"kind", "gallery_error"},
@@ -234,7 +189,7 @@ void Gallery_WebSocket_Controller::handleNewConnection(
}
try {
auto socket = std::make_shared<Gallery_WebSocket>(
connection, stream, page_session, *page_connection);
connection, stream);
connection->setContext(socket);
connection->setPingMessage(
"aethera-gallery-video", std::chrono::seconds(20));
@@ -243,8 +198,6 @@ void Gallery_WebSocket_Controller::handleNewConnection(
catch (...) {
if (const auto socket = connection->getContext<Gallery_WebSocket>())
socket->close();
else
page_session->detach(*page_connection);
try {
connection->send(nlohmann::json{
{"kind", "gallery_error"},
+2 -7
View File
@@ -1,6 +1,5 @@
#pragma once
#include "Gallery_Video_Stream.hpp"
#include "Page_Session.hpp"
#include <drogon/WebSocketController.h>
#include <exception>
#include <functional>
@@ -12,9 +11,7 @@ class Gallery_WebSocket final
: public std::enable_shared_from_this<Gallery_WebSocket> {
public:
Gallery_WebSocket(drogon::WebSocketConnectionPtr connection,
std::shared_ptr<Gallery_Video_Stream> stream,
std::shared_ptr<Page_Session> page_session,
Page_Session::Connection_Id page_connection);
std::shared_ptr<Gallery_Video_Stream> stream);
~Gallery_WebSocket();
Gallery_WebSocket(const Gallery_WebSocket&) = delete;
Gallery_WebSocket& operator=(const Gallery_WebSocket&) = delete;
@@ -34,8 +31,7 @@ public:
using Stream_Resolver =
std::function<std::shared_ptr<Gallery_Video_Stream>(std::string_view)>;
explicit Gallery_WebSocket_Controller(
Stream_Resolver resolve_stream,
std::shared_ptr<Page_Session> page_session);
Stream_Resolver resolve_stream);
void handleNewMessage(const drogon::WebSocketConnectionPtr& connection,
std::string&& message,
const drogon::WebSocketMessageType& type) override;
@@ -45,6 +41,5 @@ public:
static void initPathRouting();
private:
Stream_Resolver resolve_stream; /* 媒体组标识到共享视频流的解析入口。 */
std::shared_ptr<Page_Session> page_session; /* 所有 Gallery 与 Plot 连接共享的页面权威。 */
};
}
+14 -69
View File
@@ -4,7 +4,6 @@
#include <algorithm>
#include <atomic>
#include <chrono>
#include <mutex>
#include <stdexcept>
#include <string>
@@ -26,24 +25,17 @@ std::uint32_t input_dimension(std::uint32_t value, std::uint32_t minimum,
struct Graph_WebSocket::Private {
std::weak_ptr<drogon::WebSocketConnection> connection; /* 仅在连接存活时投递控制消息。 */
std::shared_ptr<Plot> plot; /* 该控制连接绑定的引擎与界面桥接对象。 */
std::shared_ptr<Page_Session> page_session; /* 当前连接所属的唯一页面权威。 */
Page_Session::Connection_Id page_connection{}; /* 在页面权威内登记的连接标识。 */
Plot::Stream_Id stream{}; /* Plot 完成帧通知订阅标识。 */
std::atomic_bool attached{}; /* start 成功后为真,并保证 close 只执行一次。 */
std::mutex viewport_mutex; /* 输入解码读取视口尺寸的短临界区。 */
std::uint32_t width{320}; /* 当前图集槽位对应的输入坐标宽度。 */
std::uint32_t height{192}; /* 当前图集槽位对应的输入坐标高度。 */
/* 宽高属于一个不可拆分的输入坐标版本;单个 64 位原子值是唯一状态源。 */
std::atomic_uint64_t viewport{(std::uint64_t{320} << 32U) | 192U};
};
Graph_WebSocket::Graph_WebSocket(drogon::WebSocketConnectionPtr connection,
std::shared_ptr<Plot> plot,
std::shared_ptr<Page_Session> page_session,
Page_Session::Connection_Id page_connection)
std::shared_ptr<Plot> plot)
: d(std::make_unique<Private>()) {
d->connection = std::move(connection);
d->plot = std::move(plot);
d->page_session = std::move(page_session);
d->page_connection = page_connection;
}
Graph_WebSocket::~Graph_WebSocket() { close(); }
@@ -85,9 +77,9 @@ void Graph_WebSocket::receive(std::string_view message) {
width = input_dimension(viewport->value("width", 320U), 160U, 1920U);
height = input_dimension(viewport->value("height", 192U), 120U, 1080U);
}
std::lock_guard lock(d->viewport_mutex);
d->width = width;
d->height = height;
d->viewport.store((static_cast<std::uint64_t>(width) << 32U) |
static_cast<std::uint64_t>(height),
std::memory_order_release);
return;
}
if (kind == "manual_render") {
@@ -101,13 +93,9 @@ void Graph_WebSocket::receive(std::string_view message) {
const auto type = magic_enum::enum_cast<Event_Type>(
input->value("type", std::string{}));
if (!type) return;
std::uint32_t width{};
std::uint32_t height{};
{
std::lock_guard lock(d->viewport_mutex);
width = d->width;
height = d->height;
}
const auto viewport = d->viewport.load(std::memory_order_acquire);
const auto width = static_cast<std::uint32_t>(viewport >> 32U);
const auto height = static_cast<std::uint32_t>(viewport);
Plot_Input_Event decoded;
decoded.type = *type;
const auto read_point = [&](std::string_view key,
@@ -152,20 +140,14 @@ void Graph_WebSocket::receive(std::string_view message) {
void Graph_WebSocket::close() noexcept {
if (!d->attached.exchange(false, std::memory_order_acq_rel)) return;
try { d->page_session->detach(d->page_connection); }
catch (...) {}
try { d->plot->unsubscribe(d->stream); }
catch (...) {}
}
Graph_WebSocket_Controller::Graph_WebSocket_Controller(
Plot_Resolver resolver,
std::shared_ptr<Page_Session> value_page_session)
: resolve_plot(std::move(resolver)),
page_session(std::move(value_page_session)) {
if (!resolve_plot || !page_session)
throw std::invalid_argument(
"plot WebSocket requires plot and page-session services");
Graph_WebSocket_Controller::Graph_WebSocket_Controller(Plot_Resolver resolver)
: resolve_plot(std::move(resolver)) {
if (!resolve_plot)
throw std::invalid_argument("plot WebSocket requires a plot resolver");
}
void Graph_WebSocket_Controller::initPathRouting() {
@@ -177,43 +159,8 @@ void Graph_WebSocket_Controller::handleNewConnection(
const drogon::HttpRequestPtr& request,
const drogon::WebSocketConnectionPtr& connection) {
try {
const auto page_connection = page_session->attach(
request->getParameter("session"), [weak = std::weak_ptr(connection)] {
const auto active = weak.lock();
if (!active || !active->connected()) return;
try {
active->send(nlohmann::json{
{"kind", "page_session_replaced"},
{"protocol", "aethera.page.session"},
{"version", 1}}.dump(),
drogon::WebSocketMessageType::Text);
}
catch (...) {}
try {
active->shutdown(drogon::CloseCode::kViolation,
"Page session replaced");
}
catch (...) {}
});
if (!page_connection) {
try {
connection->send(nlohmann::json{
{"kind", "page_session_rejected"},
{"protocol", "aethera.page.session"},
{"version", 1}}.dump(),
drogon::WebSocketMessageType::Text);
}
catch (...) {}
try {
connection->shutdown(drogon::CloseCode::kViolation,
"Stale page session");
}
catch (...) {}
return;
}
auto plot = resolve_plot(graph_id_from_path(request->path()));
if (!plot) {
page_session->detach(*page_connection);
try {
connection->shutdown(drogon::CloseCode::kViolation,
"Unknown Aethera plot");
@@ -223,7 +170,7 @@ void Graph_WebSocket_Controller::handleNewConnection(
}
try {
auto socket = std::make_shared<Graph_WebSocket>(
connection, std::move(plot), page_session, *page_connection);
connection, std::move(plot));
connection->setContext(socket);
connection->setPingMessage(
"aethera-gallery-plot", std::chrono::seconds(20));
@@ -232,8 +179,6 @@ void Graph_WebSocket_Controller::handleNewConnection(
catch (...) {
if (const auto socket = connection->getContext<Graph_WebSocket>())
socket->close();
else
page_session->detach(*page_connection);
try {
connection->shutdown(drogon::CloseCode::kViolation,
"Aethera plot setup failed");
+2 -8
View File
@@ -1,5 +1,4 @@
#pragma once
#include "Page_Session.hpp"
#include "Plot.hpp"
#include <drogon/WebSocketController.h>
#include <functional>
@@ -11,9 +10,7 @@ namespace aethera::web {
class Graph_WebSocket final : public std::enable_shared_from_this<Graph_WebSocket> {
public:
Graph_WebSocket(drogon::WebSocketConnectionPtr connection,
std::shared_ptr<Plot> plot,
std::shared_ptr<Page_Session> page_session,
Page_Session::Connection_Id page_connection);
std::shared_ptr<Plot> plot);
~Graph_WebSocket();
Graph_WebSocket(const Graph_WebSocket&) = delete;
Graph_WebSocket& operator=(const Graph_WebSocket&) = delete;
@@ -30,9 +27,7 @@ class Graph_WebSocket_Controller final
: public drogon::WebSocketController<Graph_WebSocket_Controller, false> {
public:
using Plot_Resolver = std::function<std::shared_ptr<Plot>(std::string_view)>;
explicit Graph_WebSocket_Controller(
Plot_Resolver resolver,
std::shared_ptr<Page_Session> page_session);
explicit Graph_WebSocket_Controller(Plot_Resolver resolver);
void handleNewMessage(const drogon::WebSocketConnectionPtr& connection,
std::string&& message,
const drogon::WebSocketMessageType& type) override;
@@ -42,6 +37,5 @@ public:
static void initPathRouting();
private:
Plot_Resolver resolve_plot; /* 图标识到引擎 Plot 的解析入口。 */
std::shared_ptr<Page_Session> page_session; /* 所有 Gallery 与 Plot 连接共享的页面权威。 */
};
}
-80
View File
@@ -1,80 +0,0 @@
#include "Page_Session.hpp"
#include <array>
#include <iomanip>
#include <mutex>
#include <random>
#include <sstream>
#include <stdexcept>
#include <unordered_map>
#include <utility>
#include <vector>
namespace aethera::web {
namespace {
std::string make_page_token() {
std::random_device random;
std::array<std::uint32_t, 4> words{};
for (auto& word : words) word = random();
std::ostringstream token;
token << std::hex << std::setfill('0');
for (const auto word : words) token << std::setw(8) << word;
return std::move(token).str();
}
}
struct Page_Session::Private {
std::mutex mutex; /* 权威页面切换与连接登记的唯一临界区。 */
std::string active_token; /* 当前唯一可建立连接的页面令牌。 */
Connection_Id next_connection{1}; /* 进程生命周期内递增的连接登记标识。 */
std::unordered_map<Connection_Id, Close_Connection>
connections; /* 当前页面拥有的连接关闭动作。 */
};
Page_Session::Page_Session() : d(std::make_unique<Private>()) {}
Page_Session::~Page_Session() = default;
std::string Page_Session::begin_page() {
auto token = make_page_token();
std::vector<Close_Connection> displaced;
{
std::lock_guard lock(d->mutex);
while (token == d->active_token) token = make_page_token();
displaced.reserve(d->connections.size());
for (auto& [connection, close_connection] : d->connections) {
static_cast<void>(connection);
displaced.push_back(std::move(close_connection));
}
d->connections.clear();
d->active_token = token;
}
/* 页面接管只在临界区内切换权威令牌。关闭第三方连接必须在锁外执行,
* 因而 WebSocket 关闭回调可以安全地反向调用 detach(),且不会阻塞新页登记。 */
for (auto& close_connection : displaced) {
try { close_connection(); }
catch (...) {
/* 单条旧连接已经失去权威资格;其关闭失败不能妨碍其余连接被隔离。 */
}
}
return token;
}
std::expected<Page_Session::Connection_Id, Page_Session::Attach_Result>
Page_Session::attach(std::string_view token,
Close_Connection close_connection) {
if (!close_connection)
throw std::invalid_argument("page connection requires a close action");
std::lock_guard lock(d->mutex);
if (token.empty() || token != d->active_token)
return std::unexpected(Attach_Result::stale_session);
const auto connection = d->next_connection++;
d->connections.emplace(connection, std::move(close_connection));
return connection;
}
void Page_Session::detach(Connection_Id connection) {
std::lock_guard lock(d->mutex);
d->connections.erase(connection);
}
}
-33
View File
@@ -1,33 +0,0 @@
#pragma once
#include <cstdint>
#include <expected>
#include <functional>
#include <memory>
#include <string>
#include <string_view>
namespace aethera::web {
class Page_Session final {
public:
using Connection_Id = std::uint64_t;
using Close_Connection = std::function<void()>;
enum class Attach_Result : std::uint8_t {
stale_session
};
Page_Session();
~Page_Session();
Page_Session(const Page_Session&) = delete;
Page_Session& operator=(const Page_Session&) = delete;
[[nodiscard]] std::string begin_page();
[[nodiscard]] std::expected<Connection_Id, Attach_Result> attach(
std::string_view token, Close_Connection close_connection);
void detach(Connection_Id connection);
private:
struct Private;
std::unique_ptr<Private> d;
};
}
+217 -176
View File
@@ -16,9 +16,7 @@
#include <initializer_list>
#include <limits>
#include <memory>
#include <mutex>
#include <optional>
#include <shared_mutex>
#include <span>
#include <stdexcept>
#include <unordered_map>
@@ -64,16 +62,15 @@ struct Frame_Pacing_Properties {
class Frame_Policy final {
public:
[[nodiscard]] Frame_Pacing_Properties snapshot() const {
std::lock_guard lock(mutex);
return pacing;
[[nodiscard]] Frame_Pacing_Properties read() const {
return *pacing.load(std::memory_order_acquire);
}
[[nodiscard]] nlohmann::json schema() const;
[[nodiscard]] nlohmann::json write_prop(std::string_view key,
const nlohmann::json& value);
private:
mutable std::mutex mutex;
Frame_Pacing_Properties pacing{};
std::atomic<std::shared_ptr<const Frame_Pacing_Properties>> pacing{
std::make_shared<const Frame_Pacing_Properties>()};
};
std::string_view pacing_mode_name(Frame_Pacing_Mode mode) {
@@ -101,22 +98,22 @@ std::string_view pixel_format_name(render_3d::Pixel_Format format) {
}
nlohmann::json Frame_Policy::schema() const {
std::lock_guard lock(mutex);
const auto current = read();
return {
{"id", "frame-analysis"}, {"label", "渲染与媒体流水线"}, {"kind", "analysis"},
{"fields", nlohmann::json::array({
{{"key", "render_enabled"}, {"label", "持续渲染与采样"}, {"editor", "boolean"},
{"editable", true}, {"description", "控制共享帧时钟是否继续调用当前 Scene;画面隐藏不会修改此项。"},
{"technical_description", "Authoritative server-side render and sampling switch."},
{"value", pacing.render_enabled}},
{"value", current.render_enabled}},
{{"key", "video_enabled"}, {"label", "图集视频传输"}, {"editor", "boolean"},
{"editable", true}, {"description", "控制完成帧是否进入页面级 RGBA 图集;默认开启,用于完整链路压测。"},
{"technical_description", "Authoritative tile publication switch for the shared gallery video."},
{"value", pacing.video_enabled}},
{"value", current.video_enabled}},
{{"key", "pacing_mode"}, {"label", "服务端帧策略"}, {"editor", "select"},
{"editable", true}, {"description", "只控制 Scene::render 的调用节奏;完成回调只负责归还帧并发布结果。"},
{"technical_description", "Render policy driven by the common 100 Hz gallery clock."},
{"value", pacing_mode_name(pacing.mode)},
{"value", pacing_mode_name(current.mode)},
{"options", nlohmann::json::array({
{{"value", "manual"}, {"label", "手动渲染"}},
{{"value", "fixed_rate"}, {"label", "固定频率"}},
@@ -126,19 +123,33 @@ nlohmann::json Frame_Policy::schema() const {
{"editable", true}, {"minimum", 0.1}, {"maximum", 100.0}, {"step", 0.1},
{"description", "固定频率模式下从页面级 100 Hz 时钟采样当前 Scene 的次数。"},
{"technical_description", "Per-plot rate selected from the common gallery render timeline."},
{"value", pacing.fixed_rate_fps}}
{"value", current.fixed_rate_fps}}
})}
};
}
nlohmann::json Frame_Policy::write_prop(std::string_view key,
const nlohmann::json& value) {
std::lock_guard lock(mutex);
const auto update = [this](auto&& edit) {
auto current = pacing.load(std::memory_order_acquire);
for (;;) {
auto next = std::make_shared<Frame_Pacing_Properties>(*current);
edit(*next);
std::shared_ptr<const Frame_Pacing_Properties> desired = next;
if (pacing.compare_exchange_weak(
current, desired, std::memory_order_release,
std::memory_order_acquire))
return desired;
}
};
if (key == "render_enabled" || key == "video_enabled") {
if (!value.is_boolean())
return {{"success", false}, {"error", "frame policy switch requires a boolean"}};
bool& target = key == "render_enabled" ? pacing.render_enabled : pacing.video_enabled;
target = value.get<bool>();
const bool target = value.get<bool>();
update([&](Frame_Pacing_Properties& next) {
(key == "render_enabled" ? next.render_enabled : next.video_enabled) =
target;
});
return {{"success", true}, {"component", "frame-analysis"}, {"key", key},
{"value", target}};
}
@@ -148,9 +159,9 @@ nlohmann::json Frame_Policy::write_prop(std::string_view key,
const auto parsed = parse_pacing_mode(value.get_ref<const std::string&>());
if (!parsed)
return {{"success", false}, {"error", "unknown frame pacing mode"}};
pacing.mode = *parsed;
update([&](Frame_Pacing_Properties& next) { next.mode = *parsed; });
return {{"success", true}, {"component", "frame-analysis"}, {"key", key},
{"value", pacing_mode_name(pacing.mode)}};
{"value", pacing_mode_name(*parsed)}};
}
if (key == "fixed_rate_fps") {
if (!value.is_number())
@@ -158,9 +169,11 @@ nlohmann::json Frame_Policy::write_prop(std::string_view key,
const double next = value.get<double>();
if (!std::isfinite(next) || next < 0.1 || next > 100.0)
return {{"success", false}, {"error", "fixed_rate_fps must be between 0.1 and 100"}};
pacing.fixed_rate_fps = next;
update([&](Frame_Pacing_Properties& properties) {
properties.fixed_rate_fps = next;
});
return {{"success", true}, {"component", "frame-analysis"}, {"key", key},
{"value", pacing.fixed_rate_fps}};
{"value", next}};
}
return {{"success", false}, {"error", "unknown frame runtime property"}};
}
@@ -332,7 +345,7 @@ struct Plot::Private {
struct Managed_Frame {
std::chrono::microseconds presentation_time{}; /* 共享页面时钟产生的媒体时间戳。 */
Frame frame{}; /* 三缓冲物理槽拥有且反复承载逻辑帧。 */
Frame_State state{Frame_State::available}; /* 本槽唯一生命周期状态。 */
std::atomic<Frame_State> state{Frame_State::available}; /* 本槽唯一生命周期状态。 */
};
struct Consumer {
@@ -340,29 +353,28 @@ struct Plot::Private {
std::uint32_t width{};
std::uint32_t height{};
};
using Consumer_Map = std::unordered_map<Stream_Id, Consumer>;
struct Stream_Snapshot {
std::vector<std::pair<Stream_Id, Consumer>> consumers;
std::shared_ptr<const Consumer_Map> consumers;
std::uint32_t width{};
std::uint32_t height{};
};
std::unique_ptr<Scene_View> view;
std::once_flag start_once;
mutable std::mutex consumers_mutex;
std::unordered_map<Stream_Id, Consumer> consumers; /* 媒体图集和诊断连接的唯一订阅表。 */
std::atomic<std::shared_ptr<const Consumer_Map>> consumers{
std::make_shared<const Consumer_Map>()}; /* 低频订阅修改发布不可变版本。 */
std::atomic_uint64_t next_stream_id{1};
std::atomic<std::shared_ptr<const std::string>> terminal_failure{}; /* 首次 Plot Unknown Failure 的唯一终止状态。 */
std::uint64_t next_frame_sequence{1};
Frame_Policy frame_policy{};
mutable std::mutex frame_mutex;
static constexpr std::size_t scene_frame_capacity{3};
std::array<Managed_Frame, scene_frame_capacity> frame_slots{}; /* Scene 借用的稳定三缓冲物理帧。 */
std::optional<std::size_t> callback_retired_slot{}; /* 上一个回调返回后才可重新开始的槽。 */
std::atomic_size_t callback_retired_slot{scene_frame_capacity};
Scene scene; /* 析构顺序保证 Scene 先停止,再释放物理帧。 */
mutable std::mutex tick_mutex;
std::optional<Plot_Render_Tick> pending_tick{}; /* 时钟拥塞时只保留尚未处理的最新时间点。 */
bool tick_task_scheduled{}; /* Taskflow 中是否已有唯一 tick 消费任务。 */
std::atomic<std::shared_ptr<const Plot_Render_Tick>> pending_tick{};
std::atomic_bool tick_task_scheduled{}; /* 唯一短任务准入;不占用 Worker 等待。 */
double last_clock_render_time_ms{-std::numeric_limits<double>::infinity()};
std::chrono::steady_clock::time_point clock_origin{std::chrono::steady_clock::now()};
std::atomic_uint64_t received_tick_count{}; /* 页面时钟交付给本 Plot 的 tick 总数。 */
@@ -373,9 +385,12 @@ struct Plot::Private {
std::atomic_uint64_t scene_rejection_count{}; /* Scene 单帧准入拒绝的提交次数。 */
std::atomic_uint64_t submitted_frame_count{}; /* 成功提交给 Scene 的帧总数。 */
std::atomic_size_t taskflow_trace_remaining{}; /* 尚待标记的实际渲染帧数。 */
mutable std::mutex taskflow_trace_mutex{}; /* 只保护低频请求结果的交换。 */
std::size_t taskflow_trace_requested_count{}; /* 当前批次请求总帧数。 */
std::vector<Taskflow_Frame_Trace> taskflow_traces{}; /* 已完成帧直接发布的 DAG 与 Observer 结果。 */
static constexpr std::size_t maximum_taskflow_trace_frames{120};
/* 高 32 位 requested,低 32 位 captured。每槽只发布一次不可变 Trace,
* GET 直接读取已发布槽位,不复制或重排整个历史容器。 */
std::atomic_uint64_t taskflow_trace_control{};
std::array<std::atomic<std::shared_ptr<const Taskflow_Frame_Trace>>,
maximum_taskflow_trace_frames> taskflow_trace_slots{};
Task_Node completion_tail{}; /* Scene 完成图中当前最后一个业务阶段。 */
std::vector<std::unique_ptr<Task_Graph>> completion_extensions{}; /* 生命周期覆盖 Scene 对子图的借用。 */
@@ -444,10 +459,9 @@ nlohmann::json Plot::Private::schema() const {
Plot::Private::Stream_Snapshot Plot::Private::stream_snapshot() const {
Stream_Snapshot result;
std::lock_guard lock(consumers_mutex);
result.consumers.reserve(consumers.size());
for (const auto& [id, consumer] : consumers) {
result.consumers.emplace_back(id, consumer);
result.consumers = consumers.load(std::memory_order_acquire);
for (const auto& [id, consumer] : *result.consumers) {
static_cast<void>(id);
if (consumer.width == 0 || consumer.height == 0) continue;
result.width = std::max(result.width, consumer.width);
result.height = std::max(result.height, consumer.height);
@@ -463,7 +477,7 @@ void Plot::Private::publish(
try {
const auto snapshot = stream_snapshot();
std::vector<Stream_Id> failed_consumers;
for (const auto& [id, consumer] : snapshot.consumers) {
for (const auto& [id, consumer] : *snapshot.consumers) {
if (!consumer.handler) continue;
try {
consumer.handler(frame);
@@ -473,27 +487,27 @@ void Plot::Private::publish(
}
}
if (failed_consumers.empty()) return;
std::lock_guard lock(consumers_mutex);
for (const auto id : failed_consumers) consumers.erase(id);
auto current = consumers.load(std::memory_order_acquire);
for (;;) {
auto next = std::make_shared<Consumer_Map>(*current);
for (const auto id : failed_consumers) next->erase(id);
std::shared_ptr<const Consumer_Map> desired = next;
if (consumers.compare_exchange_weak(
current, desired, std::memory_order_release,
std::memory_order_acquire))
break;
}
}
catch (...) {}
}
void Plot::Private::consume_tick(std::weak_ptr<Plot> lifetime) {
std::optional<Plot_Render_Tick> tick;
{
std::lock_guard lock(tick_mutex);
tick = std::exchange(pending_tick, {});
}
const auto tick = pending_tick.exchange({}, std::memory_order_acq_rel);
if (tick) clock_tick(*tick);
bool schedule_again{};
{
std::lock_guard lock(tick_mutex);
schedule_again = pending_tick.has_value();
if (!schedule_again) tick_task_scheduled = false;
}
if (schedule_again) {
tick_task_scheduled.store(false, std::memory_order_release);
if (pending_tick.load(std::memory_order_acquire) &&
!tick_task_scheduled.exchange(true, std::memory_order_acq_rel)) {
aethera::schedule_task("web.plot.tick.consume", [lifetime] {
const auto plot = lifetime.lock();
if (!plot) return;
@@ -505,7 +519,7 @@ void Plot::Private::consume_tick(std::weak_ptr<Plot> lifetime) {
void Plot::Private::clock_tick(const Plot_Render_Tick& tick) {
if (terminal_failure.load(std::memory_order_acquire)) return;
const auto pacing = frame_policy.snapshot();
const auto pacing = frame_policy.read();
if (!pacing.render_enabled || pacing.mode == Frame_Pacing_Mode::manual) {
policy_skip_count.fetch_add(1, std::memory_order_relaxed);
return;
@@ -539,46 +553,43 @@ bool Plot::Private::mark_taskflow_trace(Render_Frame& frame) {
void Plot::Private::render_frame(Plot_Render_Tick tick) {
if (terminal_failure.load(std::memory_order_acquire)) return;
const auto streams = stream_snapshot();
const auto pacing = frame_policy.snapshot();
if (!pacing.render_enabled || streams.consumers.empty()) return;
const auto pacing = frame_policy.read();
if (!pacing.render_enabled || streams.consumers->empty()) return;
std::size_t slot_index{};
Managed_Frame* managed{};
{
std::lock_guard lock(frame_mutex);
/*
* Frame_State 是 Plot 数据准备生命周期的唯一权威来源。必须在修改
* Scene/Visual 输入之前取得准入;Scene::render() 内部再拒绝已经太晚,
* 因为上一帧的异步 prepare 可能正在读取同一份业务数据。
*/
if (std::ranges::any_of(frame_slots, [](const Managed_Frame& slot) {
return slot.state == Frame_State::in_flight;
})) {
preparation_busy_count.fetch_add(1, std::memory_order_relaxed);
return;
}
const auto available = std::ranges::find_if(
frame_slots, [](const Managed_Frame& slot) {
return slot.state == Frame_State::available;
});
if (available == frame_slots.end()) {
frame_slot_busy_count.fetch_add(1, std::memory_order_relaxed);
return;
}
slot_index = static_cast<std::size_t>(
std::distance(frame_slots.begin(), available));
managed = &*available;
managed->state = Frame_State::in_flight;
managed->presentation_time =
std::chrono::duration_cast<std::chrono::microseconds>(
std::chrono::duration<double, std::milli>(
tick.time_milliseconds));
/* Frame_State 是 Plot 数据准备生命周期的唯一权威来源。槽位通过 CAS
* 准入,完成图和时钟任务不再共享一把外围锁。 */
if (std::ranges::any_of(frame_slots, [](const Managed_Frame& slot) {
return slot.state.load(std::memory_order_acquire) ==
Frame_State::in_flight;
})) {
preparation_busy_count.fetch_add(1, std::memory_order_relaxed);
return;
}
for (std::size_t index = 0; index < frame_slots.size(); ++index) {
auto expected = Frame_State::available;
if (!frame_slots[index].state.compare_exchange_strong(
expected, Frame_State::in_flight,
std::memory_order_acq_rel, std::memory_order_acquire))
continue;
slot_index = index;
managed = &frame_slots[index];
break;
}
if (!managed) {
frame_slot_busy_count.fetch_add(1, std::memory_order_relaxed);
return;
}
managed->presentation_time =
std::chrono::duration_cast<std::chrono::microseconds>(
std::chrono::duration<double, std::milli>(tick.time_milliseconds));
const auto rollback_unsubmitted = [this, slot_index] {
std::lock_guard lock(frame_mutex);
auto& slot = frame_slots[slot_index];
if (slot.state == Frame_State::in_flight)
slot.state = Frame_State::available;
auto expected = Frame_State::in_flight;
static_cast<void>(slot.state.compare_exchange_strong(
expected, Frame_State::available, std::memory_order_acq_rel,
std::memory_order_acquire));
};
bool taskflow_trace_claimed{};
const auto restore_taskflow_trace_claim = [this, &taskflow_trace_claimed] {
@@ -662,36 +673,32 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) {
void Plot::Private::publish_completed_frame() {
Render_Frame* frame{};
Managed_Frame* managed{};
{
std::lock_guard lock(frame_mutex);
/*
* Scene 在用户回调返回之后才标记 callback_finished,因此当前回调
* 不能立即重启自己的物理帧。下一个串行完成回调开始时,上一个
* callback_retired 槽已经确定离开 Scene,可安全归还。三槽由此在
* 不改变 Scene 回调契约的前提下形成稳定的循环生命周期。
*/
if (callback_retired_slot) {
auto& retired = frame_slots[*callback_retired_slot];
if (retired.state != Frame_State::callback_retired)
throw std::logic_error("Plot retired frame state is inconsistent");
retired.state = Frame_State::available;
callback_retired_slot.reset();
}
for (std::size_t index = 0; index < frame_slots.size(); ++index) {
if (frame_slots[index].state != Frame_State::in_flight) continue;
if (managed)
throw std::logic_error("Plot has multiple active Scene frames");
frame = std::visit(
[](const auto& value) -> Render_Frame* { return value.get(); },
frame_slots[index].frame);
managed = &frame_slots[index];
}
if (!managed || managed->state != Frame_State::in_flight)
throw std::logic_error("Scene completion graph has no active Plot frame");
/* Scene 在用户回调返回之后才标记 callback_finished。下一次串行完成图
* 开始时,上一 callback_retired 槽才可归还。 */
const auto retired_index = callback_retired_slot.exchange(
scene_frame_capacity, std::memory_order_acq_rel);
if (retired_index != scene_frame_capacity) {
auto expected = Frame_State::callback_retired;
if (!frame_slots[retired_index].state.compare_exchange_strong(
expected, Frame_State::available, std::memory_order_acq_rel,
std::memory_order_acquire))
throw std::logic_error("Plot retired frame state is inconsistent");
}
for (std::size_t index = 0; index < frame_slots.size(); ++index) {
if (frame_slots[index].state.load(std::memory_order_acquire) !=
Frame_State::in_flight) continue;
if (managed)
throw std::logic_error("Plot has multiple active Scene frames");
frame = std::visit(
[](const auto& value) -> Render_Frame* { return value.get(); },
frame_slots[index].frame);
managed = &frame_slots[index];
}
if (!managed)
throw std::logic_error("Scene completion graph has no active Plot frame");
try {
const auto pacing = frame_policy.snapshot();
const auto pacing = frame_policy.read();
const auto identity = frame->identity();
Frame_Identity rendered_identity = identity;
std::shared_ptr<const std::vector<std::byte>> pixel_storage;
@@ -748,30 +755,44 @@ void Plot::Private::retire_completed_frame(Render_Frame* frame) {
if (!frame)
throw std::invalid_argument("Plot received a null completed frame");
std::size_t slot_index{scene_frame_capacity};
{
std::lock_guard lock(frame_mutex);
for (std::size_t index = 0; index < frame_slots.size(); ++index) {
auto* address = std::visit(
[](const auto& value) -> Render_Frame* { return value.get(); },
frame_slots[index].frame);
if (address != frame) continue;
if (frame_slots[index].state != Frame_State::in_flight)
throw std::logic_error("completed Plot frame is not in flight");
slot_index = index;
break;
}
if (slot_index == scene_frame_capacity)
throw std::logic_error(
"frame callback has no externally owned active frame");
if (callback_retired_slot)
throw std::logic_error("Plot has more than one callback-retired frame");
frame_slots[slot_index].state = Frame_State::callback_retired;
callback_retired_slot = slot_index;
for (std::size_t index = 0; index < frame_slots.size(); ++index) {
auto* address = std::visit(
[](const auto& value) -> Render_Frame* { return value.get(); },
frame_slots[index].frame);
if (address != frame) continue;
auto expected = Frame_State::in_flight;
if (!frame_slots[index].state.compare_exchange_strong(
expected, Frame_State::callback_retired,
std::memory_order_acq_rel, std::memory_order_acquire))
throw std::logic_error("completed Plot frame is not in flight");
slot_index = index;
break;
}
if (slot_index == scene_frame_capacity)
throw std::logic_error(
"frame callback has no externally owned active frame");
auto no_retired = scene_frame_capacity;
if (!callback_retired_slot.compare_exchange_strong(
no_retired, slot_index, std::memory_order_acq_rel,
std::memory_order_acquire))
throw std::logic_error("Plot has more than one callback-retired frame");
if (frame->taskflow_trace_requested()) {
auto trace = frame->taskflow_trace();
std::lock_guard lock(taskflow_trace_mutex);
taskflow_traces.push_back(std::move(trace));
auto trace = std::make_shared<const Taskflow_Frame_Trace>(
frame->taskflow_trace());
auto control = 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 (captured >= requested) break;
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(
control, next, std::memory_order_release,
std::memory_order_acquire))
break;
}
}
}
@@ -837,14 +858,18 @@ Plot::Stream_Id Plot::subscribe(Stream_Handler handler) {
throw std::invalid_argument("Plot subscription requires a handler");
ensure_started();
const auto id = d->next_stream_id.fetch_add(1, std::memory_order_relaxed);
Stream_Handler notification;
{
std::lock_guard lock(d->consumers_mutex);
d->consumers.emplace(id, Private::Consumer{std::move(handler)});
if (d->terminal_failure.load(std::memory_order_acquire))
notification = d->consumers.at(id).handler;
const auto notification = handler;
auto current = d->consumers.load(std::memory_order_acquire);
for (;;) {
auto next = std::make_shared<Private::Consumer_Map>(*current);
next->emplace(id, Private::Consumer{handler});
std::shared_ptr<const Private::Consumer_Map> desired = next;
if (d->consumers.compare_exchange_weak(
current, desired, std::memory_order_release,
std::memory_order_acquire))
break;
}
if (notification) {
if (d->terminal_failure.load(std::memory_order_acquire)) {
const auto failure = d->terminal_failure.load(std::memory_order_acquire);
try {
notification(std::make_shared<const Plot_Stream_Frame>(
@@ -863,35 +888,44 @@ Plot::Stream_Id Plot::subscribe(Stream_Handler handler) {
}
void Plot::unsubscribe(Stream_Id stream) {
std::lock_guard lock(d->consumers_mutex);
d->consumers.erase(stream);
auto current = d->consumers.load(std::memory_order_acquire);
while (current->contains(stream)) {
auto next = std::make_shared<Private::Consumer_Map>(*current);
next->erase(stream);
std::shared_ptr<const Private::Consumer_Map> desired = next;
if (d->consumers.compare_exchange_weak(
current, desired, std::memory_order_release,
std::memory_order_acquire))
return;
}
}
void Plot::configure_stream(Stream_Id stream, std::uint32_t width,
std::uint32_t height) {
std::lock_guard lock(d->consumers_mutex);
const auto found = d->consumers.find(stream);
if (found == d->consumers.end()) return;
found->second.width = width;
found->second.height = height;
auto current = d->consumers.load(std::memory_order_acquire);
for (;;) {
const auto found = current->find(stream);
if (found == current->end()) return;
auto next = std::make_shared<Private::Consumer_Map>(*current);
auto& consumer = next->at(stream);
consumer.width = width;
consumer.height = height;
std::shared_ptr<const Private::Consumer_Map> desired = next;
if (d->consumers.compare_exchange_weak(
current, desired, std::memory_order_release,
std::memory_order_acquire))
return;
}
}
void Plot::schedule_render(Plot_Render_Tick tick) {
ensure_started();
if (d->terminal_failure.load(std::memory_order_acquire)) return;
d->received_tick_count.fetch_add(1, std::memory_order_relaxed);
bool schedule{};
{
std::lock_guard lock(d->tick_mutex);
if (d->pending_tick)
d->coalesced_tick_count.fetch_add(1, std::memory_order_relaxed);
d->pending_tick = tick;
if (!d->tick_task_scheduled) {
d->tick_task_scheduled = true;
schedule = true;
}
}
if (!schedule) return;
const auto next = std::make_shared<const Plot_Render_Tick>(std::move(tick));
if (d->pending_tick.exchange(next, std::memory_order_acq_rel))
d->coalesced_tick_count.fetch_add(1, std::memory_order_relaxed);
if (d->tick_task_scheduled.exchange(true, std::memory_order_acq_rel)) return;
const auto weak = weak_from_this();
aethera::schedule_task("web.plot.tick.consume", [weak] {
const auto owner = weak.lock();
@@ -995,7 +1029,7 @@ nlohmann::json Plot::diagnostics() const {
}
}, d->scene);
const auto pacing = d->frame_policy.snapshot();
const auto pacing = d->frame_policy.read();
const auto stream = d->stream_snapshot();
nlohmann::json supported_formats = nlohmann::json::array();
if (is_3d) {
@@ -1068,7 +1102,6 @@ nlohmann::json Plot::diagnostics() const {
{"callback_max_ms", milliseconds(gpu.callback_max_ns)},
{"callback_failure_count", gpu.callback_failure_count},
{"backpressure_count", gpu.backpressure_count},
{"backpressure_wait_ms", milliseconds(gpu.backpressure_wait_ns)},
{"fault_count", gpu.fault_count},
{"abandoned_count", gpu.abandoned_count},
{"stopping", gpu.stopping}};
@@ -1079,29 +1112,37 @@ nlohmann::json Plot::diagnostics() const {
}
void Plot::request_taskflow_trace(std::size_t frame_count) {
if (frame_count == 0 || frame_count > 120)
if (frame_count == 0 ||
frame_count > Private::maximum_taskflow_trace_frames)
throw std::invalid_argument("Taskflow trace frame_count must be between 1 and 120");
ensure_started();
{
std::lock_guard lock(d->taskflow_trace_mutex);
if (d->taskflow_trace_requested_count != d->taskflow_traces.size())
auto control = d->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 Taskflow frame trace request is already active");
d->taskflow_traces.clear();
d->taskflow_traces.reserve(frame_count);
d->taskflow_trace_requested_count = frame_count;
const auto next = static_cast<std::uint64_t>(frame_count) << 32U;
if (d->taskflow_trace_control.compare_exchange_weak(
control, next, std::memory_order_release,
std::memory_order_acquire))
break;
}
for (auto& slot : d->taskflow_trace_slots)
slot.store({}, std::memory_order_release);
d->taskflow_trace_remaining.store(frame_count, std::memory_order_release);
}
nlohmann::json Plot::taskflow_trace() const {
nlohmann::json frames = nlohmann::json::array();
std::size_t requested{};
{
std::lock_guard lock(d->taskflow_trace_mutex);
requested = d->taskflow_trace_requested_count;
for (const auto& trace : d->taskflow_traces)
frames.push_back(taskflow_trace_json(trace));
}
const auto control = d->taskflow_trace_control.load(
std::memory_order_acquire);
const auto requested = static_cast<std::uint32_t>(control >> 32U);
const auto captured = static_cast<std::uint32_t>(control);
for (std::uint32_t index = 0; index < captured; ++index)
if (const auto trace = d->taskflow_trace_slots[index].load(
std::memory_order_acquire))
frames.push_back(taskflow_trace_json(*trace));
const auto remaining = d->taskflow_trace_remaining.load(
std::memory_order_acquire);
return {
+50 -126
View File
@@ -1,6 +1,5 @@
#include "WebRtc_Video_Session.hpp"
#include <concurrentqueue-1.0.5/blockingconcurrentqueue.h>
#include <double_buffer/mechanism.hpp>
#include <nlohmann/json.hpp>
#include <rtc/rtc.hpp>
#include <atomic>
@@ -37,32 +36,6 @@ void update_maximum(std::atomic_uint64_t& maximum,
}
struct WebRtc_Video_Session::Private {
struct State_Tag {};
struct State : double_buffer::State_Type<State_Tag> {
std::size_t outstanding_video_frames{};
std::size_t transport_buffered_bytes{};
std::uint64_t queued_frame_count{};
std::uint64_t rejected_frame_count{};
std::uint64_t queued_byte_count{};
std::uint64_t sent_frame_count{};
std::uint64_t sent_byte_count{};
std::uint64_t send_failure_count{};
std::uint64_t send_total_ns{};
std::uint64_t send_max_ns{};
std::uint64_t current_send_ns{};
std::uint64_t media_lock_wait_count{};
std::uint64_t media_lock_wait_total_ns{};
std::uint64_t media_lock_wait_max_ns{};
std::uint64_t readiness_check_count{};
std::uint64_t readiness_reject_count{};
std::uint64_t close_join_total_ns{};
std::uint64_t close_join_max_ns{};
std::uint64_t current_close_join_ns{};
bool track_open{};
bool closed{};
bool sender_running{};
bool operator==(const State&) const = default;
};
struct Callback_State {
Signal_Handler signal_handler; /* SDP/ICE 信令的唯一交付出口。 */
Ready_Handler ready_handler; /* Track 打开后请求首个 IDR。 */
@@ -82,9 +55,11 @@ struct WebRtc_Video_Session::Private {
std::shared_ptr<Callback_State> callbacks;
std::shared_ptr<rtc::PeerConnection> peer;
std::shared_ptr<rtc::Track> video_track;
std::atomic<std::shared_ptr<rtc::Track>> video_track{};
std::shared_ptr<rtc::RtpPacketizationConfig> rtp_config;
std::mutex media_mutex; /* 发送、信令修改与关闭的生命周期边界。 */
/* libdatachannel 没有为同一 Peer 的 SDP/ICE 修改与 close 提供可并发契约;
* 这里只串行化第三方对象生命周期。视频热路径通过原子 Track 发布,不取锁。 */
std::mutex media_mutex;
moodycamel::BlockingConcurrentQueue<Encoded_Video_Frame>
pending_frames{4}; /* 尚未发送的最新编码帧入口。 */
std::thread sender_thread{}; /* 唯一的 RTP/SRTP 发送资源域。 */
@@ -100,9 +75,6 @@ struct WebRtc_Video_Session::Private {
std::atomic_uint64_t send_started_ns{};
std::atomic_bool send_active{};
std::atomic_bool sender_loop_active{};
std::atomic_uint64_t media_lock_wait_count{};
std::atomic_uint64_t media_lock_wait_total_ns{};
std::atomic_uint64_t media_lock_wait_max_ns{};
std::atomic_size_t transport_buffered_bytes{};
std::atomic_uint64_t readiness_check_count{};
std::atomic_uint64_t readiness_reject_count{};
@@ -110,18 +82,6 @@ struct WebRtc_Video_Session::Private {
std::atomic_uint64_t close_join_total_ns{};
std::atomic_uint64_t close_join_max_ns{};
std::atomic_bool close_join_active{};
/* WebRTC 是 web_server 会话域,不借用 Scene State。热路径只写原子计数;
* 低频 diagnostics 请求在这里聚合并交换会话自己的发布双缓冲。 */
mutable std::mutex state_publication_mutex;
mutable double_buffer::Publish_Double_Buffer<State> state{};
void record_media_lock_wait(std::uint64_t started) noexcept {
const auto elapsed = steady_nanoseconds() - started;
media_lock_wait_count.fetch_add(1, std::memory_order_relaxed);
media_lock_wait_total_ns.fetch_add(elapsed, std::memory_order_relaxed);
update_maximum(media_lock_wait_max_ns, elapsed);
}
Private(Signal_Handler signal_handler, Ready_Handler ready_handler,
Failure_Handler failure_handler)
: callbacks(std::make_shared<Callback_State>()) {
@@ -146,20 +106,14 @@ struct WebRtc_Video_Session::Private {
callbacks->closed.load(std::memory_order_acquire))
return;
try {
const auto lock_started = steady_nanoseconds();
std::shared_ptr<rtc::Track> track;
{
std::lock_guard lock(media_mutex);
record_media_lock_wait(lock_started);
if (!video_track ||
!callbacks->track_open.load(std::memory_order_acquire) ||
callbacks->closed.load(std::memory_order_acquire)) {
outstanding_video_frames.fetch_sub(
1, std::memory_order_release);
rejected_frame_count.fetch_add(1, std::memory_order_relaxed);
continue;
}
track = video_track;
const auto track = video_track.load(std::memory_order_acquire);
if (!track ||
!callbacks->track_open.load(std::memory_order_acquire) ||
callbacks->closed.load(std::memory_order_acquire)) {
outstanding_video_frames.fetch_sub(
1, std::memory_order_release);
rejected_frame_count.fetch_add(1, std::memory_order_relaxed);
continue;
}
/* rtc::binary 与编码器 access unit 使用相同的 byte vector。
* 这里只在锁内取得 Track 的共享所有权;packetizer 和网络发送均为
@@ -298,7 +252,7 @@ void WebRtc_Video_Session::start() {
peer->setLocalDescription();
d->peer = std::move(peer);
d->video_track = std::move(video_track);
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(); });
}
@@ -325,7 +279,7 @@ WebRtc_Video_Session::Send_Result WebRtc_Video_Session::send(
if (frame.annex_b.empty() ||
!d->callbacks->track_open.load(std::memory_order_acquire) ||
d->callbacks->closed.load(std::memory_order_acquire) ||
!d->sender_thread.joinable()) {
!d->sender_loop_active.load(std::memory_order_acquire)) {
d->rejected_frame_count.fetch_add(1, std::memory_order_relaxed);
return Send_Result::not_open;
}
@@ -362,18 +316,12 @@ WebRtc_Video_Session::Send_Result WebRtc_Video_Session::send(
bool WebRtc_Video_Session::can_accept_video() const noexcept {
d->readiness_check_count.fetch_add(1, std::memory_order_relaxed);
const auto lock_started = steady_nanoseconds();
std::shared_ptr<rtc::Track> track;
{
std::lock_guard lock(d->media_mutex);
d->record_media_lock_wait(lock_started);
track = d->video_track;
}
const auto track = d->video_track.load(std::memory_order_acquire);
const auto buffered = track ? track->bufferedAmount() : 0;
d->transport_buffered_bytes.store(buffered, std::memory_order_relaxed);
const bool ready = d->callbacks->track_open.load(std::memory_order_acquire) &&
!d->callbacks->closed.load(std::memory_order_acquire) &&
d->sender_thread.joinable() &&
d->sender_loop_active.load(std::memory_order_acquire) &&
track &&
buffered < maximum_transport_buffered_bytes &&
d->outstanding_video_frames.load(std::memory_order_acquire) <
@@ -384,38 +332,18 @@ bool WebRtc_Video_Session::can_accept_video() const noexcept {
}
nlohmann::json WebRtc_Video_Session::diagnostics() const {
std::lock_guard publication_lock(d->state_publication_mutex);
auto& state_to_publish = *d->state.current;
const auto now = steady_nanoseconds();
state_to_publish.outstanding_video_frames =
const auto outstanding_video_frames =
d->outstanding_video_frames.load(std::memory_order_relaxed);
state_to_publish.transport_buffered_bytes =
d->transport_buffered_bytes.load(std::memory_order_relaxed);
state_to_publish.queued_frame_count = d->queued_frame_count.load(std::memory_order_relaxed);
state_to_publish.rejected_frame_count = d->rejected_frame_count.load(std::memory_order_relaxed);
state_to_publish.queued_byte_count = d->queued_byte_count.load(std::memory_order_relaxed);
state_to_publish.sent_frame_count = d->sent_frame_count.load(std::memory_order_relaxed);
state_to_publish.sent_byte_count = d->sent_byte_count.load(std::memory_order_relaxed);
state_to_publish.send_failure_count = d->send_failure_count.load(std::memory_order_relaxed);
state_to_publish.send_total_ns = d->send_total_ns.load(std::memory_order_relaxed);
state_to_publish.send_max_ns = d->send_max_ns.load(std::memory_order_relaxed);
state_to_publish.current_send_ns = d->send_active.load(std::memory_order_acquire)
? now - d->send_started_ns.load(std::memory_order_relaxed) : 0;
state_to_publish.media_lock_wait_count = d->media_lock_wait_count.load(std::memory_order_relaxed);
state_to_publish.media_lock_wait_total_ns = d->media_lock_wait_total_ns.load(std::memory_order_relaxed);
state_to_publish.media_lock_wait_max_ns = d->media_lock_wait_max_ns.load(std::memory_order_relaxed);
state_to_publish.readiness_check_count = d->readiness_check_count.load(std::memory_order_relaxed);
state_to_publish.readiness_reject_count = d->readiness_reject_count.load(std::memory_order_relaxed);
state_to_publish.close_join_total_ns = d->close_join_total_ns.load(std::memory_order_relaxed);
state_to_publish.close_join_max_ns = d->close_join_max_ns.load(std::memory_order_relaxed);
state_to_publish.current_close_join_ns = d->close_join_active.load(std::memory_order_acquire)
? now - d->close_join_started_ns.load(std::memory_order_relaxed) : 0;
state_to_publish.track_open = d->callbacks->track_open.load(std::memory_order_relaxed);
state_to_publish.closed = d->callbacks->closed.load(std::memory_order_relaxed);
state_to_publish.sender_running = d->sender_loop_active.load(std::memory_order_relaxed);
d->state.advance();
const auto& published = *d->state.pending;
const auto queued_frame_count =
d->queued_frame_count.load(std::memory_order_relaxed);
const auto sent_frame_count =
d->sent_frame_count.load(std::memory_order_relaxed);
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;
};
@@ -423,34 +351,30 @@ nlohmann::json WebRtc_Video_Session::diagnostics() const {
{"kind", "webrtc_transport_state"},
{"protocol", "aethera.gallery.webrtc"},
{"version", 1},
{"track_open", published.track_open},
{"closed", published.closed},
{"sender_running", published.sender_running},
{"outstanding_video_frames", published.outstanding_video_frames},
{"transport_buffered_bytes", published.transport_buffered_bytes},
{"queued_frame_count", published.queued_frame_count},
{"rejected_frame_count", published.rejected_frame_count},
{"queued_megabytes", static_cast<double>(published.queued_byte_count) /
{"track_open", d->callbacks->track_open.load(std::memory_order_relaxed)},
{"closed", d->callbacks->closed.load(std::memory_order_relaxed)},
{"sender_running", d->sender_loop_active.load(std::memory_order_relaxed)},
{"outstanding_video_frames", outstanding_video_frames},
{"transport_buffered_bytes", d->transport_buffered_bytes.load(std::memory_order_relaxed)},
{"queued_frame_count", queued_frame_count},
{"rejected_frame_count", d->rejected_frame_count.load(std::memory_order_relaxed)},
{"queued_megabytes", static_cast<double>(d->queued_byte_count.load(std::memory_order_relaxed)) /
(1024.0 * 1024.0)},
{"sent_frame_count", published.sent_frame_count},
{"sent_megabytes", static_cast<double>(published.sent_byte_count) /
{"sent_frame_count", sent_frame_count},
{"sent_megabytes", static_cast<double>(d->sent_byte_count.load(std::memory_order_relaxed)) /
(1024.0 * 1024.0)},
{"send_failure_count", published.send_failure_count},
{"send_average_ms", published.sent_frame_count == 0 ? 0.0 :
milliseconds(published.send_total_ns) /
static_cast<double>(published.sent_frame_count)},
{"send_max_ms", milliseconds(published.send_max_ns)},
{"current_send_ms", milliseconds(published.current_send_ns)},
{"media_lock_wait_count", published.media_lock_wait_count},
{"media_lock_wait_average_ms", published.media_lock_wait_count == 0 ? 0.0 :
milliseconds(published.media_lock_wait_total_ns) /
static_cast<double>(published.media_lock_wait_count)},
{"media_lock_wait_max_ms", milliseconds(published.media_lock_wait_max_ns)},
{"readiness_check_count", published.readiness_check_count},
{"readiness_reject_count", published.readiness_reject_count},
{"close_join_total_ms", milliseconds(published.close_join_total_ns)},
{"close_join_max_ms", milliseconds(published.close_join_max_ns)},
{"current_close_join_ms", milliseconds(published.current_close_join_ns)}};
{"send_failure_count", d->send_failure_count.load(std::memory_order_relaxed)},
{"send_average_ms", sent_frame_count == 0 ? 0.0 :
milliseconds(send_total_ns) / static_cast<double>(sent_frame_count)},
{"send_max_ms", milliseconds(d->send_max_ns.load(std::memory_order_relaxed))},
{"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)}};
}
void WebRtc_Video_Session::close() noexcept {
@@ -477,11 +401,11 @@ void WebRtc_Video_Session::close() noexcept {
d->outstanding_video_frames.store(0, std::memory_order_release);
std::lock_guard lock(d->media_mutex);
try { if (d->video_track) d->video_track->close(); }
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->video_track.reset();
d->rtp_config.reset();
d->peer.reset();
}
+2 -14
View File
@@ -3,7 +3,6 @@
#include "Gallery_WebSocket.hpp"
#include "Graph_WebSocket.hpp"
#include "Gallery_Plots.hpp"
#include "Page_Session.hpp"
#include <drogon/drogon.h>
#include <nlohmann/json.hpp>
#include <algorithm>
@@ -176,26 +175,15 @@ int run_web_server(std::uint16_t port, const std::filesystem::path& asset_root)
};
add_media_groups("2d", std::move(gallery_2d));
add_media_groups("3d", std::move(gallery_3d));
auto page_session = std::make_shared<Page_Session>();
auto resolve_plot = [plots](std::string_view id) { return find_plot(*plots, id); };
auto websocket = std::make_shared<Graph_WebSocket_Controller>(
resolve_plot, page_session);
auto websocket = std::make_shared<Graph_WebSocket_Controller>(resolve_plot);
auto gallery_websocket = std::make_shared<Gallery_WebSocket_Controller>(
[gallery_streams](std::string_view id) {
const auto found = gallery_streams->find(std::string{id});
return found == gallery_streams->end() ? nullptr : found->second;
}, page_session);
});
auto& app = drogon::app();
app.registerHandler("/page/session", [page_session](
const drogon::HttpRequestPtr&,
std::function<void(const drogon::HttpResponsePtr&)>&& callback) {
callback(json_response({
{"protocol", "aethera.page.session"},
{"version", 1},
{"token", page_session->begin_page()}}));
}, {drogon::Post});
app.registerHandler("/plot", [plots, plot_media](const drogon::HttpRequestPtr&,
std::function<void(const drogon::HttpResponsePtr&)>&& callback) {
nlohmann::json result = nlohmann::json::array();
+42 -38
View File
@@ -3,19 +3,23 @@
#include <atomic>
#include <cstring>
#include <limits>
#include <mutex>
#include <stdexcept>
#include <utility>
namespace aethera::web::detail {
struct Gallery_Frame_Atlas::Private {
struct Source {
mutable std::mutex mutex; /* 仅保护当前 Plot 的最近完成快照。 */
std::shared_ptr<const Plot_Pixel_Frame> latest_completion; /* 最近逻辑完成身份;像素可为空。 */
std::shared_ptr<const Plot_Pixel_Frame> latest_pixels; /* 最近可合成的实际 RGBA 画面。 */
std::uint64_t composited_rendered_sequence{}; /* 图集像素当前包含的真实画面序号。 */
std::uint64_t completion_count{}; /* 当前槽位接受的逻辑完成回调总数。 */
std::uint64_t rendered_frame_count{}; /* 当前槽位接受的不同真实画面总数。 */
struct Published {
std::shared_ptr<const Plot_Pixel_Frame> latest_completion;
std::shared_ptr<const Plot_Pixel_Frame> latest_pixels;
std::uint64_t completion_count{};
std::uint64_t rendered_frame_count{};
};
/* Plot 完成线程发布不可变版本,图集线程只读取同一个版本;不再把四个
* */
std::atomic<std::shared_ptr<const Published>> published{
std::make_shared<const Published>()};
std::atomic_uint64_t composited_rendered_sequence{};
};
Gallery_Atlas_Description description{}; /* 尺寸与槽位布局的唯一权威描述。 */
@@ -83,7 +87,6 @@ Gallery_Frame_Atlas::Accept_Frame_Result Gallery_Frame_Atlas::accept_frame(
return Accept_Frame_Result::invalid_frame;
}
auto& source = *d->sources[slot];
std::lock_guard lock(source.mutex);
if (frame->pixels) {
const auto candidate = static_cast<int>(frame->layout);
int unknown{-1};
@@ -95,21 +98,29 @@ Gallery_Frame_Atlas::Accept_Frame_Result Gallery_Frame_Atlas::accept_frame(
return Accept_Frame_Result::invalid_frame;
}
}
if (source.latest_completion &&
frame->sequence <= source.latest_completion->sequence) {
d->rejected_frame_count.fetch_add(1, std::memory_order_relaxed);
return Accept_Frame_Result::stale_frame;
auto current = source.published.load(std::memory_order_acquire);
for (;;) {
if (current->latest_completion &&
frame->sequence <= current->latest_completion->sequence) {
d->rejected_frame_count.fetch_add(1, std::memory_order_relaxed);
return Accept_Frame_Result::stale_frame;
}
auto next = std::make_shared<Private::Source::Published>(*current);
if (frame->rendered_sequence != 0 &&
(!current->latest_completion ||
frame->rendered_sequence !=
current->latest_completion->rendered_sequence))
++next->rendered_frame_count;
next->latest_completion = frame;
if (frame->pixels && frame->rendered_sequence != 0)
next->latest_pixels = frame;
++next->completion_count;
std::shared_ptr<const Private::Source::Published> desired = next;
if (source.published.compare_exchange_weak(
current, std::move(desired), std::memory_order_release,
std::memory_order_acquire))
return Accept_Frame_Result::accepted;
}
if (frame->rendered_sequence != 0 &&
(!source.latest_completion ||
frame->rendered_sequence !=
source.latest_completion->rendered_sequence))
++source.rendered_frame_count;
source.latest_completion = frame;
if (frame->pixels && frame->rendered_sequence != 0)
source.latest_pixels = std::move(frame);
++source.completion_count;
return Accept_Frame_Result::accepted;
}
Gallery_Atlas_Composition Gallery_Frame_Atlas::compose() {
@@ -126,17 +137,9 @@ Gallery_Atlas_Composition Gallery_Frame_Atlas::compose() {
static_cast<std::size_t>(d->description.tile_width) * 4U;
for (std::size_t slot = 0; slot < d->sources.size(); ++slot) {
auto& source = *d->sources[slot];
std::shared_ptr<const Plot_Pixel_Frame> latest_completion;
std::shared_ptr<const Plot_Pixel_Frame> latest_pixels;
std::uint64_t completion_count{};
std::uint64_t rendered_frame_count{};
{
std::lock_guard lock(source.mutex);
latest_completion = source.latest_completion;
latest_pixels = source.latest_pixels;
completion_count = source.completion_count;
rendered_frame_count = source.rendered_frame_count;
}
const auto published = source.published.load(std::memory_order_acquire);
const auto& latest_completion = published->latest_completion;
const auto& latest_pixels = published->latest_pixels;
if (!latest_completion) {
++result.missing_tile_count;
result.sources.push_back({});
@@ -146,7 +149,8 @@ Gallery_Atlas_Composition Gallery_Frame_Atlas::compose() {
++result.missing_tile_count;
}
else if (latest_pixels->rendered_sequence !=
source.composited_rendered_sequence) {
source.composited_rendered_sequence.load(
std::memory_order_acquire)) {
const auto& layout = d->description.sources[slot];
for (std::uint32_t y = 0; y < d->description.tile_height; ++y) {
const auto source_offset =
@@ -160,8 +164,8 @@ Gallery_Atlas_Composition Gallery_Frame_Atlas::compose() {
latest_pixels->pixels->data() + source_offset,
source_row_bytes);
}
source.composited_rendered_sequence =
latest_pixels->rendered_sequence;
source.composited_rendered_sequence.store(
latest_pixels->rendered_sequence, std::memory_order_release);
++result.fresh_tile_count;
}
result.sources.push_back({
@@ -169,8 +173,8 @@ Gallery_Atlas_Composition Gallery_Frame_Atlas::compose() {
latest_completion->correlation_id,
latest_completion->rendered_sequence,
latest_completion->rendered_correlation_id,
completion_count,
rendered_frame_count});
published->completion_count,
published->rendered_frame_count});
}
result.pixels = d->pixels;
return result;
+27 -13
View File
@@ -13,10 +13,13 @@ struct Gallery_Frame_Clock::Private {
std::chrono::nanoseconds interval{};
Tick_Handler tick_handler{};
Failure_Handler failure_handler{};
/* condition_variable_any 的 stop_token 等待和 jthread 替换必须共享生命周期
* */
std::mutex lifecycle_mutex{};
std::condition_variable_any wake{};
std::jthread thread{};
std::atomic_bool running{};
std::atomic_uint64_t generation{}; /* 区分 stop 期间接管的新时钟。 */
Private(double frame_rate, Tick_Handler value_tick_handler,
Failure_Handler value_failure_handler)
@@ -32,14 +35,16 @@ struct Gallery_Frame_Clock::Private {
throw std::invalid_argument("gallery frame clock requires a tick handler");
}
void report(std::exception_ptr failure) noexcept {
running.store(false, std::memory_order_release);
void report(std::uint64_t active_generation,
std::exception_ptr failure) noexcept {
if (generation.load(std::memory_order_acquire) == active_generation)
running.store(false, std::memory_order_release);
if (!failure_handler) return;
try { failure_handler(std::move(failure)); }
catch (...) {}
}
void run(std::stop_token stop) noexcept {
void run(std::stop_token stop, std::uint64_t active_generation) noexcept {
const auto origin = std::chrono::steady_clock::now();
std::uint64_t sequence{};
auto deadline = origin + interval;
@@ -60,11 +65,12 @@ struct Gallery_Frame_Clock::Private {
deadline += interval * ((now - deadline) / interval + 1);
}
catch (...) {
report(std::current_exception());
report(active_generation, std::current_exception());
return;
}
}
running.store(false, std::memory_order_release);
if (generation.load(std::memory_order_acquire) == active_generation)
running.store(false, std::memory_order_release);
}
};
@@ -79,17 +85,25 @@ Gallery_Frame_Clock::~Gallery_Frame_Clock() { stop(); }
void Gallery_Frame_Clock::start() {
std::lock_guard lock(d->lifecycle_mutex);
if (d->running.exchange(true, std::memory_order_acq_rel)) return;
d->thread = std::jthread([state = d](std::stop_token stop) {
state->run(stop);
const auto generation = d->generation.fetch_add(
1, std::memory_order_acq_rel) + 1;
d->thread = std::jthread([state = d, generation](std::stop_token stop) {
state->run(stop, generation);
});
}
void Gallery_Frame_Clock::stop() noexcept {
if (!d || !d->running.exchange(false, std::memory_order_acq_rel)) return;
d->thread.request_stop();
d->wake.notify_all();
if (d->thread.joinable() &&
d->thread.get_id() != std::this_thread::get_id())
d->thread.join();
if (!d) return;
std::jthread thread;
{
std::lock_guard lock(d->lifecycle_mutex);
if (!d->running.exchange(false, std::memory_order_acq_rel)) return;
d->thread.request_stop();
d->wake.notify_all();
thread = std::move(d->thread);
}
if (!thread.joinable()) return;
if (thread.get_id() == std::this_thread::get_id()) thread.detach();
else thread.join();
}
}
-36
View File
@@ -1,36 +0,0 @@
#include "web_server/src/Page_Session.hpp"
#include <gtest/gtest.h>
namespace aethera::web {
TEST(Page_Session, New_Page_Displaces_Previous_Page_Without_Waiting_For_Close) {
Page_Session session;
const auto first_token = session.begin_page();
std::size_t displaced{};
const auto first_connection = session.attach(
first_token, [&displaced] { ++displaced; });
ASSERT_TRUE(first_connection.has_value());
const auto second_token = session.begin_page();
EXPECT_NE(first_token, second_token);
EXPECT_EQ(displaced, 1U);
EXPECT_EQ(session.attach(first_token, [] {}).error(),
Page_Session::Attach_Result::stale_session);
const auto second_connection = session.attach(second_token, [] {});
EXPECT_TRUE(second_connection.has_value());
session.detach(*first_connection);
session.detach(*second_connection);
}
TEST(Page_Session, One_Page_Can_Attach_All_Of_Its_Connections) {
Page_Session session;
const auto token = session.begin_page();
const auto gallery = session.attach(token, [] {});
const auto plot = session.attach(token, [] {});
ASSERT_TRUE(gallery.has_value());
ASSERT_TRUE(plot.has_value());
EXPECT_NE(*gallery, *plot);
session.detach(*gallery);
session.detach(*plot);
}
}
@@ -1,5 +1,6 @@
#include <frame_statistics.hpp>
#include <gtest/gtest.h>
#include <cmath>
#include <limits>
namespace aethera::web {
@@ -52,4 +53,18 @@ TEST(Sliding_Statistics, Estimates_Quantiles_Without_Growing_The_Window) {
EXPECT_NEAR(state.p95, 950.0, 15.0);
EXPECT_NEAR(state.p99, 990.0, 15.0);
}
TEST(Sliding_Statistics, Keeps_Approximate_Quantiles_Ordered_For_Nonstationary_Input) {
Sliding_Statistics statistics{16};
for (std::size_t phase = 0; phase < 200; ++phase) {
const double baseline = phase % 2 == 0 ? 1'000'000.0 : -1'000'000.0;
for (std::size_t sample = 0; sample < 17; ++sample) {
const auto state = statistics.submit(
baseline + static_cast<double>(sample * sample));
EXPECT_TRUE(std::isfinite(state.trimmed_average));
EXPECT_LE(state.p50, state.p95);
EXPECT_LE(state.p95, state.p99);
}
}
}
}