Skip to content

Commit 8f4d694

Browse files
committed
Harden subscription finalization and catalog teardown timing
1 parent 1e53b82 commit 8f4d694

1 file changed

Lines changed: 42 additions & 36 deletions

File tree

src/transport/moqt_session.cpp

Lines changed: 42 additions & 36 deletions
Original file line numberDiff line numberDiff line change
@@ -1333,6 +1333,7 @@ TransportStatus serve_subscriptions(PublisherTransport& transport,
13331333
std::map<std::uint64_t, ActiveSubscription> active_subscriptions;
13341334
bool fin = false;
13351335
bool served_any_subscription = false;
1336+
std::optional<std::chrono::steady_clock::time_point> catalog_last_served_at;
13361337
std::uint64_t first_media_time_us = 0;
13371338
bool first_media_time_set = false;
13381339
const auto pacing_start = std::chrono::steady_clock::now();
@@ -1633,6 +1634,33 @@ TransportStatus serve_subscriptions(PublisherTransport& transport,
16331634
}
16341635

16351636
if (!active_subscriptions.empty()) {
1637+
std::vector<std::uint64_t> completed_request_ids_to_finalize;
1638+
for (const auto& [request_id, active] : active_subscriptions) {
1639+
if (active.completed) {
1640+
completed_request_ids_to_finalize.push_back(request_id);
1641+
}
1642+
}
1643+
for (const auto request_id : completed_request_ids_to_finalize) {
1644+
const auto active_it = active_subscriptions.find(request_id);
1645+
if (active_it == active_subscriptions.end()) {
1646+
continue;
1647+
}
1648+
const TransportStatus finalize_status =
1649+
finalize_subscription(transport,
1650+
control_stream_id,
1651+
request_id,
1652+
active_it->second.sender->stream_count(),
1653+
completed_request_ids);
1654+
if (!finalize_status.ok) {
1655+
return finalize_status;
1656+
}
1657+
active_subscriptions.erase(active_it);
1658+
}
1659+
if (!completed_request_ids_to_finalize.empty()) {
1660+
served_any_subscription = true;
1661+
continue;
1662+
}
1663+
16361664
std::vector<std::uint8_t> chunk;
16371665
bool immediate_fin = false;
16381666
const TransportStatus read_status =
@@ -1726,6 +1754,9 @@ TransportStatus serve_subscriptions(PublisherTransport& transport,
17261754
<< " track=" << object.track_name
17271755
<< " group=" << object.group_id << " object=" << object.object_id
17281756
<< " bytes=" << object_payload_size(source_object) << '\n';
1757+
if (object.track_name == "catalog") {
1758+
catalog_last_served_at = std::chrono::steady_clock::now();
1759+
}
17291760
trace_csv_write_served("served",
17301761
send_seq,
17311762
trace_elapsed_ms(std::chrono::steady_clock::now()),
@@ -1748,32 +1779,6 @@ TransportStatus serve_subscriptions(PublisherTransport& transport,
17481779
continue;
17491780
}
17501781

1751-
std::vector<std::uint64_t> completed_request_ids_to_finalize;
1752-
for (const auto& [request_id, active] : active_subscriptions) {
1753-
if (active.completed) {
1754-
completed_request_ids_to_finalize.push_back(request_id);
1755-
}
1756-
}
1757-
for (const auto request_id : completed_request_ids_to_finalize) {
1758-
const auto active_it = active_subscriptions.find(request_id);
1759-
if (active_it == active_subscriptions.end()) {
1760-
continue;
1761-
}
1762-
const TransportStatus finalize_status =
1763-
finalize_subscription(transport,
1764-
control_stream_id,
1765-
request_id,
1766-
active_it->second.sender->stream_count(),
1767-
completed_request_ids);
1768-
if (!finalize_status.ok) {
1769-
return finalize_status;
1770-
}
1771-
active_subscriptions.erase(active_it);
1772-
}
1773-
if (!completed_request_ids_to_finalize.empty()) {
1774-
served_any_subscription = true;
1775-
continue;
1776-
}
17771782
}
17781783

17791784
if (fin) {
@@ -1811,6 +1816,17 @@ TransportStatus serve_subscriptions(PublisherTransport& transport,
18111816
return transport.write_stream(control_stream_id, encode_publish_namespace_done_message(namespace_message), false);
18121817
}
18131818

1819+
if (send_namespace_done && catalog_last_served_at.has_value()) {
1820+
const auto now = std::chrono::steady_clock::now();
1821+
const auto elapsed = std::chrono::duration_cast<std::chrono::milliseconds>(now - *catalog_last_served_at);
1822+
const auto desired_grace = std::min(subscriber_timeout, std::chrono::milliseconds(250));
1823+
if (elapsed < desired_grace) {
1824+
const auto remaining = desired_grace - elapsed;
1825+
std::cerr << "[moqt-session] grace wait after catalog publish: " << remaining.count() << " ms\n";
1826+
std::this_thread::sleep_for(remaining);
1827+
}
1828+
}
1829+
18141830
pending_control_bytes = std::move(buffer);
18151831
if (!send_namespace_done) {
18161832
return TransportStatus::success();
@@ -2531,16 +2547,6 @@ TransportStatus MoqtSession::publish_live(std::istream& input,
25312547
return TransportStatus::success();
25322548
};
25332549

2534-
if (auto_forward_) {
2535-
const auto catalog_alias_it = alias_by_track.find("catalog");
2536-
if (catalog_alias_it != alias_by_track.end()) {
2537-
status = send_catalog(catalog_alias_it->second);
2538-
if (!status.ok) {
2539-
return status;
2540-
}
2541-
}
2542-
}
2543-
25442550
// Phase 4: Stream media from stdin.
25452551
// Use a reader thread so we can also handle control messages.
25462552
struct LiveMediaQueue {

0 commit comments

Comments
 (0)