Skip to content

Commit 67edd5c

Browse files
committed
Handle relay subscribe flow for forward=0
1 parent fc6c45f commit 67edd5c

3 files changed

Lines changed: 59 additions & 62 deletions

File tree

README.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -196,7 +196,7 @@ Current status as of March 13, 2026:
196196
- QUIC handshake succeeds against `draft-14.cloudflare.mediaoverquic.com:443`, `interop-relay.cloudflare.mediaoverquic.com:443`, and `moq-relay.red5.net:8443`
197197
- `CLIENT_SETUP` succeeds and the client prints the negotiated connection ID to stdout after setup
198198
- `PUBLISH_NAMESPACE` is accepted with `PUBLISH_NAMESPACE_OK`
199-
- with `--forward 0`, the current client then waits for inbound `SUBSCRIBE_NAMESPACE` / `SUBSCRIBE`
199+
- with `--forward 0`, the current client waits for inbound `SUBSCRIBE`; relays may consume `SUBSCRIBE_NAMESPACE` themselves and only forward `SUBSCRIBE` to the publisher
200200
- the Cloudflare endpoints accepted setup and namespace announce in testing, but did not issue subscriptions, so the publish attempt timed out waiting for control-stream data
201201
- with `--forward 1`, `moq-relay.red5.net:8443` now progresses through `PUBLISH_OK` for the catalog and media tracks, after which the client begins sending object streams
202202
- `fb.mvfst.net:9448` now accepts the draft-14 publish flow end-to-end after switching `PUBLISH`, `PUBLISH_OK`, and `PUBLISH_ERROR` control messages to `u16` outer lengths; the current draft-16 flow is still rejected with MOQT application error `3` (`PROTOCOL_VIOLATION`) immediately after setup

src/transport/moqt_session.cpp

Lines changed: 24 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -16,10 +16,6 @@ namespace openmoq::publisher::transport {
1616

1717
namespace {
1818

19-
bool control_message_complete(std::span<const std::uint8_t> bytes, std::size_t& message_size) {
20-
return next_control_message(bytes, message_size);
21-
}
22-
2319
bool trace_enabled() {
2420
static const bool enabled = std::getenv("OPENMOQ_PICOQUIC_TRACE") != nullptr;
2521
return enabled;
@@ -419,33 +415,31 @@ TransportStatus serve_subscriptions(PublisherTransport& transport,
419415
std::uint64_t first_media_time_us = 0;
420416
bool first_media_time_set = false;
421417
const auto pacing_start = std::chrono::steady_clock::now();
422-
if (subscribe.forward != 0) {
423-
for (const auto& object : plan.objects) {
424-
if (object.track_name != subscribe.track_name || !object_matches_filter(object, subscribe)) {
425-
continue;
426-
}
427-
const auto payload = object_payload(object);
428-
if (payload.empty()) {
429-
return TransportStatus::failure("transport publish requires materialized object payloads");
430-
}
431-
if (object.kind == openmoq::publisher::CmsfObjectKind::kMedia && !first_media_time_set) {
432-
first_media_time_us = object.media_time_us;
433-
first_media_time_set = true;
434-
}
435-
pace_until(pacing_start, first_media_time_us, object, paced);
436-
std::uint64_t stream_id = 0;
437-
write_status = transport.open_stream(StreamDirection::kUnidirectional, stream_id);
438-
if (!write_status.ok) {
439-
return write_status;
440-
}
441-
write_status = transport.write_stream(stream_id,
442-
encode_object_stream(draft, track_it->second.alias, object, payload),
443-
true);
444-
if (!write_status.ok) {
445-
return write_status;
446-
}
447-
++stream_count;
418+
for (const auto& object : plan.objects) {
419+
if (object.track_name != subscribe.track_name || !object_matches_filter(object, subscribe)) {
420+
continue;
421+
}
422+
const auto payload = object_payload(object);
423+
if (payload.empty()) {
424+
return TransportStatus::failure("transport publish requires materialized object payloads");
425+
}
426+
if (object.kind == openmoq::publisher::CmsfObjectKind::kMedia && !first_media_time_set) {
427+
first_media_time_us = object.media_time_us;
428+
first_media_time_set = true;
429+
}
430+
pace_until(pacing_start, first_media_time_us, object, paced);
431+
std::uint64_t stream_id = 0;
432+
write_status = transport.open_stream(StreamDirection::kUnidirectional, stream_id);
433+
if (!write_status.ok) {
434+
return write_status;
435+
}
436+
write_status = transport.write_stream(stream_id,
437+
encode_object_stream(draft, track_it->second.alias, object, payload),
438+
true);
439+
if (!write_status.ok) {
440+
return write_status;
448441
}
442+
++stream_count;
449443
}
450444

451445
write_status = transport.write_stream(control_stream_id,

tests/moqt_session_test.cpp

Lines changed: 34 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -164,7 +164,8 @@ std::vector<std::uint8_t> encode_subscribe_namespace_message(std::uint64_t reque
164164

165165
std::vector<std::uint8_t> encode_subscribe_message(std::uint64_t request_id,
166166
std::string_view track_namespace,
167-
std::string_view track_name) {
167+
std::string_view track_name,
168+
std::uint8_t forward) {
168169
std::vector<std::uint8_t> payload = encode_varint(request_id);
169170
const std::vector<std::uint8_t> tuple_len = encode_varint(1);
170171
const std::vector<std::uint8_t> component_len = encode_varint(track_namespace.size());
@@ -176,7 +177,7 @@ std::vector<std::uint8_t> encode_subscribe_message(std::uint64_t request_id,
176177
payload.insert(payload.end(), track_name.begin(), track_name.end());
177178
payload.push_back(0x80);
178179
payload.push_back(0x01);
179-
payload.push_back(0x01);
180+
payload.push_back(forward);
180181
const std::vector<std::uint8_t> filter_type = encode_varint(0x03);
181182
const std::vector<std::uint8_t> start_group = encode_varint(0);
182183
const std::vector<std::uint8_t> start_object = encode_varint(0);
@@ -216,11 +217,15 @@ std::vector<std::uint8_t> encode_publish_ok_message(DraftVersion draft, std::uin
216217
void queue_subscribe_requests(MockTransport& transport,
217218
DraftVersion draft,
218219
std::string_view track_namespace,
219-
std::initializer_list<std::pair<std::uint64_t, std::string>> requests) {
220+
std::initializer_list<std::pair<std::uint64_t, std::string>> requests,
221+
bool include_subscribe_namespace = false,
222+
std::uint8_t forward = 0) {
220223
transport.reads[0].push_back(encode_publish_namespace_ok_message(draft, 0));
221-
transport.reads[0].push_back(encode_subscribe_namespace_message(1, track_namespace));
224+
if (include_subscribe_namespace) {
225+
transport.reads[0].push_back(encode_subscribe_namespace_message(1, track_namespace));
226+
}
222227
for (const auto& [request_id, track_name] : requests) {
223-
transport.reads[0].push_back(encode_subscribe_message(request_id, track_namespace, track_name));
228+
transport.reads[0].push_back(encode_subscribe_message(request_id, track_namespace, track_name, forward));
224229
}
225230
}
226231

@@ -408,9 +413,9 @@ int main() {
408413
materialize_publish_plan(make_span_backed_plan(DraftVersion::kDraft14), source_bytes);
409414

410415
status = session.publish(materialized);
411-
ok &= expect(status.ok, "expected publish to succeed with announce plus subscribe flow");
412-
ok &= expect(transport.writes.size() == 10,
413-
"expected setup, namespace, subscribe_namespace_ok, two subscribe_ok, two object streams, two publish_done, namespace_done");
416+
ok &= expect(status.ok, "expected publish to succeed with relay subscribe flow");
417+
ok &= expect(transport.writes.size() == 9,
418+
"expected setup, namespace, two subscribe_ok, two object streams, two publish_done, namespace_done");
414419
ok &= expect(!transport.writes[0].bytes.empty() && transport.writes[0].bytes.front() == 0x20,
415420
"expected binary CLIENT_SETUP message type");
416421
authority.clear();
@@ -425,21 +430,20 @@ int main() {
425430
ok &= expect(transport.writes[1].bytes == std::vector<std::uint8_t>({0x06, 0x00, 0x0b, 0x00, 0x01, 0x07, 0x69, 0x6e,
426431
0x74, 0x65, 0x72, 0x6f, 0x70, 0x00}),
427432
"expected namespace write to use the configured track namespace");
428-
ok &= expect(message_type(transport.writes[2].bytes) == 0x12, "expected SUBSCRIBE_NAMESPACE_OK");
429-
ok &= expect(message_type(transport.writes[3].bytes) == 0x04, "expected first SUBSCRIBE_OK");
430-
ok &= expect(transport.writes[4].stream_id == 2, "expected first object stream to be unidirectional stream 2");
431-
ok &= expect(!transport.writes[4].bytes.empty() && transport.writes[4].bytes.front() == 0x14,
433+
ok &= expect(message_type(transport.writes[2].bytes) == 0x04, "expected first SUBSCRIBE_OK");
434+
ok &= expect(transport.writes[3].stream_id == 2, "expected first object stream to be unidirectional stream 2");
435+
ok &= expect(!transport.writes[3].bytes.empty() && transport.writes[3].bytes.front() == 0x14,
432436
"expected first object stream to use a subgroup header with explicit subgroup ID");
433-
ok &= expect(transport.writes[4].fin, "expected first object stream write to set FIN");
434-
ok &= expect(message_type(transport.writes[5].bytes) == 0x0b, "expected first PUBLISH_DONE");
435-
ok &= expect(message_type(transport.writes[6].bytes) == 0x04, "expected second SUBSCRIBE_OK");
436-
ok &= expect(transport.writes[7].stream_id == 6, "expected second object stream to be unidirectional stream 6");
437-
ok &= expect(!transport.writes[7].bytes.empty() && transport.writes[7].bytes.front() == 0x14,
437+
ok &= expect(transport.writes[3].fin, "expected first object stream write to set FIN");
438+
ok &= expect(message_type(transport.writes[4].bytes) == 0x0b, "expected first PUBLISH_DONE");
439+
ok &= expect(message_type(transport.writes[5].bytes) == 0x04, "expected second SUBSCRIBE_OK");
440+
ok &= expect(transport.writes[6].stream_id == 6, "expected second object stream to be unidirectional stream 6");
441+
ok &= expect(!transport.writes[6].bytes.empty() && transport.writes[6].bytes.front() == 0x14,
438442
"expected second object stream to use a subgroup header with explicit subgroup ID");
439-
ok &= expect(transport.writes[7].fin, "expected second object stream write to set FIN");
440-
ok &= expect(message_type(transport.writes[8].bytes) == 0x0b, "expected second PUBLISH_DONE");
441-
ok &= expect(message_type(transport.writes[9].bytes) == 0x09, "expected PUBLISH_NAMESPACE_DONE");
442-
ok &= expect(transport.writes[9].bytes == std::vector<std::uint8_t>({0x09, 0x00, 0x09, 0x01, 0x07, 0x69, 0x6e,
443+
ok &= expect(transport.writes[6].fin, "expected second object stream write to set FIN");
444+
ok &= expect(message_type(transport.writes[7].bytes) == 0x0b, "expected second PUBLISH_DONE");
445+
ok &= expect(message_type(transport.writes[8].bytes) == 0x09, "expected PUBLISH_NAMESPACE_DONE");
446+
ok &= expect(transport.writes[8].bytes == std::vector<std::uint8_t>({0x09, 0x00, 0x09, 0x01, 0x07, 0x69, 0x6e,
443447
0x74, 0x65, 0x72, 0x6f, 0x70}),
444448
"expected draft-14 PUBLISH_NAMESPACE_DONE to contain the configured track namespace");
445449
}
@@ -498,7 +502,7 @@ int main() {
498502
materialize_publish_plan(make_span_backed_plan(DraftVersion::kDraft16), source_bytes);
499503
status = draft16_session.publish(draft16_materialized);
500504
ok &= expect(status.ok, "expected draft-16 publish to succeed");
501-
ok &= expect(draft16_transport.writes.size() == 10, "expected draft-16 announce/subscribe control/object sequence");
505+
ok &= expect(draft16_transport.writes.size() == 9, "expected draft-16 relay subscribe control/object sequence");
502506
ok &= expect(!draft16_transport.writes[0].bytes.empty() && draft16_transport.writes[0].bytes.front() == 0x20,
503507
"expected draft-16 binary CLIENT_SETUP message type");
504508
authority.clear();
@@ -515,16 +519,15 @@ int main() {
515519
0x69, 0x6e, 0x74, 0x65, 0x72, 0x6f,
516520
0x70, 0x00}),
517521
"expected draft-16 namespace write to use the configured track namespace");
518-
ok &= expect(message_type(draft16_transport.writes[2].bytes) == 0x07, "expected draft-16 REQUEST_OK");
519-
ok &= expect(message_type(draft16_transport.writes[3].bytes) == 0x04, "expected first draft-16 SUBSCRIBE_OK");
520-
ok &= expect(draft16_transport.writes[4].stream_id == 2, "expected first draft-16 object stream");
521-
ok &= expect(message_type(draft16_transport.writes[5].bytes) == 0x0b, "expected first draft-16 PUBLISH_DONE");
522-
ok &= expect(message_type(draft16_transport.writes[6].bytes) == 0x04, "expected second draft-16 SUBSCRIBE_OK");
523-
ok &= expect(draft16_transport.writes[7].stream_id == 6, "expected second draft-16 object stream");
524-
ok &= expect(message_type(draft16_transport.writes[8].bytes) == 0x0b, "expected second draft-16 PUBLISH_DONE");
525-
ok &= expect(message_type(draft16_transport.writes[9].bytes) == 0x09,
522+
ok &= expect(message_type(draft16_transport.writes[2].bytes) == 0x04, "expected first draft-16 SUBSCRIBE_OK");
523+
ok &= expect(draft16_transport.writes[3].stream_id == 2, "expected first draft-16 object stream");
524+
ok &= expect(message_type(draft16_transport.writes[4].bytes) == 0x0b, "expected first draft-16 PUBLISH_DONE");
525+
ok &= expect(message_type(draft16_transport.writes[5].bytes) == 0x04, "expected second draft-16 SUBSCRIBE_OK");
526+
ok &= expect(draft16_transport.writes[6].stream_id == 6, "expected second draft-16 object stream");
527+
ok &= expect(message_type(draft16_transport.writes[7].bytes) == 0x0b, "expected second draft-16 PUBLISH_DONE");
528+
ok &= expect(message_type(draft16_transport.writes[8].bytes) == 0x09,
526529
"expected draft-16 PUBLISH_NAMESPACE_DONE");
527-
ok &= expect(draft16_transport.writes[9].bytes == std::vector<std::uint8_t>({0x09, 0x00, 0x01, 0x00}),
530+
ok &= expect(draft16_transport.writes[8].bytes == std::vector<std::uint8_t>({0x09, 0x00, 0x01, 0x00}),
528531
"expected draft-16 PUBLISH_NAMESPACE_DONE to contain only the request ID");
529532

530533
const std::vector<std::uint8_t> split_server_setup = encode_server_setup_message({

0 commit comments

Comments
 (0)