Skip to content

Commit 6cd12c8

Browse files
afrindmeta-codesync[bot]
authored andcommitted
Support TrackStatus in moxygen relay (#107)
Summary: Added track_status to MoQRelay issue: #97 co-authored by akash-a-n Pull Request resolved: #107 Reviewed By: sandarsh Differential Revision: D95152978 Pulled By: afrind fbshipit-source-id: bd8296e628f78f332355e4cae192477c58d09a3a
1 parent 25c3d71 commit 6cd12c8

4 files changed

Lines changed: 178 additions & 0 deletions

File tree

moxygen/relay/MoQRelay.cpp

Lines changed: 70 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -988,6 +988,76 @@ folly::coro::Task<Publisher::FetchResult> MoQRelay::fetch(
988988
fetch, std::move(consumer), std::move(upstreamSession));
989989
}
990990

991+
folly::coro::Task<Publisher::TrackStatusResult> MoQRelay::trackStatus(
992+
TrackStatus trackStatus) {
993+
XLOG(DBG1) << __func__ << " ftn=" << trackStatus.fullTrackName;
994+
995+
if (trackStatus.fullTrackName.trackNamespace.empty()) {
996+
co_return folly::makeUnexpected(TrackStatusError(
997+
{trackStatus.requestID,
998+
TrackStatusErrorCode::TRACK_NOT_EXIST,
999+
"namespace required"}));
1000+
}
1001+
1002+
auto subscriptionIt = subscriptions_.find(trackStatus.fullTrackName);
1003+
if (subscriptionIt != subscriptions_.end() &&
1004+
subscriptionIt->second.forwarder->numForwardingSubscribers() > 0) {
1005+
// We have active subscription - answer directly from local forwarder state
1006+
auto& subscription = subscriptionIt->second;
1007+
auto& forwarder = subscription.forwarder;
1008+
1009+
TrackStatusCode statusCode = TrackStatusCode::TRACK_NOT_STARTED;
1010+
// forwarder->largest() being set means: we have actually
1011+
// received at least one object for this track.
1012+
// subscription.handle being non-null means: the relay still has a
1013+
// live upstream Publisher::SubscriptionHandle for this track
1014+
if (forwarder->largest()) {
1015+
if (subscription.handle) {
1016+
statusCode = TrackStatusCode::IN_PROGRESS;
1017+
} else {
1018+
statusCode = TrackStatusCode::UNKNOWN;
1019+
}
1020+
}
1021+
1022+
TrackStatusOk trackStatusOk;
1023+
trackStatusOk.requestID = trackStatus.requestID;
1024+
trackStatusOk.groupOrder = forwarder->groupOrder();
1025+
trackStatusOk.largest = forwarder->largest();
1026+
trackStatusOk.fullTrackName = trackStatus.fullTrackName;
1027+
trackStatusOk.statusCode = statusCode;
1028+
1029+
XLOG(DBG1) << "Returning local track status for "
1030+
<< trackStatus.fullTrackName
1031+
<< " statusCode=" << (uint32_t)statusCode;
1032+
co_return trackStatusOk;
1033+
} else {
1034+
// No subscription - forward to upstream
1035+
auto upstreamSession =
1036+
findPublishNamespaceSession(trackStatus.fullTrackName.trackNamespace);
1037+
1038+
if (!upstreamSession) {
1039+
XLOG(DBG1) << "No upstream session for track: "
1040+
<< trackStatus.fullTrackName;
1041+
co_return folly::makeUnexpected(
1042+
TrackStatusError{
1043+
trackStatus.requestID,
1044+
TrackStatusErrorCode::TRACK_NOT_EXIST,
1045+
"no such namespace or track"});
1046+
}
1047+
1048+
// Forward the trackStatus request to the upstream publisher session
1049+
auto result = co_await upstreamSession->trackStatus(trackStatus);
1050+
1051+
if (result.hasError()) {
1052+
XLOG(DBG1) << "Upstream trackStatus failed: "
1053+
<< result.error().reasonPhrase;
1054+
} else {
1055+
XLOG(DBG1) << "Upstream trackStatus succeeded";
1056+
}
1057+
co_return result;
1058+
}
1059+
}
1060+
9911061
void MoQRelay::onEmpty(MoQForwarder* forwarder) {
9921062
auto subscriptionIt = subscriptions_.find(forwarder->fullTrackName());
9931063
if (subscriptionIt == subscriptions_.end()) {

moxygen/relay/MoQRelay.h

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -57,6 +57,9 @@ class MoQRelay : public Publisher,
5757
XLOG(INFO) << "Processing goaway uri=" << goaway.newSessionUri;
5858
}
5959

60+
folly::coro::Task<Publisher::TrackStatusResult> trackStatus(
61+
TrackStatus req) override;
62+
6063
std::shared_ptr<MoQSession> findPublishNamespaceSession(
6164
const TrackNamespace& ns);
6265

moxygen/relay/test/MoQRelayTest.cpp

Lines changed: 99 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2176,4 +2176,103 @@ TEST_F(MoQRelayTest, ExactNamespaceSubscriberReceivesPublishNamespace) {
21762176
removeSession(subscriber);
21772177
}
21782178

2179+
// Test: TrackStatus on non-existent track
2180+
TEST_F(MoQRelayTest, TrackStatusNonExistentTrack) {
2181+
auto clientSession = createMockSession();
2182+
2183+
// Request trackStatus for a track that doesn't exist
2184+
TrackStatus trackStatus;
2185+
trackStatus.fullTrackName = kTestTrackName;
2186+
trackStatus.requestID = RequestID(1);
2187+
2188+
withSessionContext(clientSession, [&]() {
2189+
auto task = relay_->trackStatus(trackStatus);
2190+
auto res = folly::coro::blockingWait(std::move(task), exec_.get());
2191+
2192+
// Should return error indicating track not found
2193+
EXPECT_FALSE(res.hasValue());
2194+
EXPECT_EQ(res.error().errorCode, TrackStatusErrorCode::TRACK_NOT_EXIST);
2195+
EXPECT_FALSE(res.error().reasonPhrase.empty());
2196+
});
2197+
2198+
removeSession(clientSession);
2199+
}
2200+
2201+
// Test: TrackStatus on existing track - returns forwarder state (no upstream
2202+
// call)
2203+
TEST_F(MoQRelayTest, TrackStatusSuccessfulForward) {
2204+
auto publisherSession = createMockSession();
2205+
auto clientSession = createMockSession();
2206+
2207+
doPublish(publisherSession, kTestTrackName);
2208+
2209+
auto consumer = createMockConsumer();
2210+
subscribeToTrack(clientSession, kTestTrackName, consumer, RequestID(1));
2211+
2212+
TrackStatus trackStatus;
2213+
trackStatus.fullTrackName = kTestTrackName;
2214+
trackStatus.requestID = RequestID(2);
2215+
2216+
withSessionContext(clientSession, [&]() {
2217+
auto task = relay_->trackStatus(trackStatus);
2218+
auto res = folly::coro::blockingWait(std::move(task), exec_.get());
2219+
2220+
// Should return status from local forwarder
2221+
// Since no data was sent, statusCode should be TRACK_NOT_STARTED
2222+
EXPECT_TRUE(res.hasValue());
2223+
EXPECT_EQ(res.value().statusCode, TrackStatusCode::TRACK_NOT_STARTED);
2224+
EXPECT_EQ(res.value().fullTrackName, kTestTrackName);
2225+
});
2226+
2227+
removeSession(clientSession);
2228+
exec_->drive();
2229+
removeSession(publisherSession);
2230+
}
2231+
2232+
// Test: TrackStatus using namespace prefix matching (no exact subscription)
2233+
// Verifies that when there's no exact subscription but a publisher has
2234+
// published a matching namespace prefix, the relay correctly routes
2235+
// TRACK_STATUS upstream using prefix matching
2236+
TEST_F(MoQRelayTest, TrackStatusViaPrefixMatching) {
2237+
auto publisher = createMockSession();
2238+
auto requester = createMockSession();
2239+
2240+
// Publisher publishes namespace but NOT the specific track
2241+
doPublishNamespace(publisher, kTestNamespace);
2242+
2243+
// No exact subscription exists for kTestTrackName, so trackStatus should
2244+
// use prefix matching to find the publisher
2245+
2246+
// Mock the upstream trackStatus call
2247+
TrackStatusOk statusOk;
2248+
statusOk.requestID = RequestID(1);
2249+
statusOk.trackAlias = TrackAlias(0);
2250+
statusOk.largest = AbsoluteLocation{50, 25};
2251+
2252+
EXPECT_CALL(*publisher, trackStatus(_)).WillOnce([statusOk](auto /*ts*/) {
2253+
return folly::coro::makeTask<Publisher::TrackStatusResult>(statusOk);
2254+
});
2255+
2256+
// Execute trackStatus from requester's perspective
2257+
TrackStatus trackStatus;
2258+
trackStatus.requestID = RequestID(1);
2259+
trackStatus.fullTrackName = kTestTrackName;
2260+
2261+
withSessionContext(requester, [&]() {
2262+
auto task = relay_->trackStatus(trackStatus);
2263+
auto result = folly::coro::blockingWait(std::move(task), exec_.get());
2264+
2265+
// Should successfully forward via prefix matching and return the result
2266+
EXPECT_TRUE(result.hasValue())
2267+
<< "TrackStatus via namespace prefix matching should succeed";
2268+
EXPECT_EQ(result.value().requestID, RequestID(1));
2269+
EXPECT_TRUE(result.value().largest.has_value());
2270+
EXPECT_EQ(result.value().largest->group, 50);
2271+
EXPECT_EQ(result.value().largest->object, 25);
2272+
});
2273+
2274+
removeSession(publisher);
2275+
removeSession(requester);
2276+
}
2277+
21792278
} // namespace moxygen::test

moxygen/test/MockMoQSession.h

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -61,6 +61,12 @@ class MockMoQSession : public MoQSession {
6161
(),
6262
(const, override));
6363

64+
MOCK_METHOD(
65+
folly::coro::Task<Publisher::TrackStatusResult>,
66+
trackStatus,
67+
(TrackStatus),
68+
(override));
69+
6470
RequestID peekNextRequestID() {
6571
return RequestID(nextRequestID_++);
6672
}

0 commit comments

Comments
 (0)