Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 0 additions & 1 deletion quest/m1/test-flakes-2/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,6 @@ Public API: none. Wire: none.
## Required

- [Media late join](/quest/m1/test-flakes-2/media-late-join.md) - a late joiner shows video promptly and catches up to live, and the check asserts that under load
- [Import catalog finish](/quest/m1/test-flakes-2/import-catalog-finish.md) - `moq-cli`'s subprocess EOF catalog-finish test holds up under load with event-based fixture coordination
- [Relay restart rebind](/quest/m1/test-flakes-2/relay-restart-rebind.md) - the crash drill restarts on its original UDP address under concurrent load
- [Impaired handshake](/quest/m1/test-flakes-2/impaired-handshake.md) - the impaired cluster drills' clients never time out while connecting
- [Media audio tone](/quest/m1/media-audio-tone.md) - the `just test media` audio-tone check passes under load, fixed at its cause
22 changes: 0 additions & 22 deletions quest/m1/test-flakes-2/import-catalog-finish.md

This file was deleted.

34 changes: 22 additions & 12 deletions rs/moq-cli/tests/import.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,10 +5,23 @@ use std::process::Stdio;
use std::time::Duration;

use tokio::io::AsyncWriteExt;
use tokio::process::Child;

const TIMEOUT: Duration = Duration::from_secs(10);
const BBB: &[u8] = include_bytes!("../../moq-mux/src/container/ts/test_data/bbb.ts");

/// Await `step` while `moq` still holds stdin open, failing as soon as the process exits.
///
/// Until stdin ends, an exit is the publisher failing, so report its status rather than
/// waiting out `TIMEOUT` and blaming the step.
async fn running<T>(child: &mut Child, what: &str, step: impl Future<Output = T>) -> T {
tokio::select! {
out = step => out,
status = child.wait() => panic!("moq exited with {} before {what}", status.expect("wait for moq")),
_ = tokio::time::sleep(TIMEOUT) => panic!("{what} timed out"),
}
}

#[tokio::test]
async fn import_delivers_the_catalog_finish_at_eof() {
let _ = moq_tokio::crypto::install_default();
Expand Down Expand Up @@ -58,27 +71,24 @@ async fn import_delivers_the_catalog_finish_at_eof() {
.await
.expect("subscriber connects");

// `Live` can come first while the relay has yet to learn the publisher's route.
while !matches!(
tokio::time::timeout(TIMEOUT, announced.next())
.await
.expect("announce timed out")
.expect("origin closed"),
moq_tokio::moq_net::announce::Event::Start(_)
) {}
let broadcast = tokio::time::timeout(TIMEOUT, consumer.request_broadcast("demo", None))
let event = running(&mut child, "the announce", announced.next())
.await
.expect("origin closed");
let moq_tokio::moq_net::announce::Event::Start(announce) = event else {
panic!("the first announce event is {event:?}, not a start");
};
assert_eq!(announce.prefix.as_str(), "demo");
let broadcast = running(&mut child, "the broadcast", consumer.request_broadcast("demo", None))
.await
.expect("request timed out")
.expect("announced broadcast resolves");
let mut catalogs = hang::catalog::Catalog::<()>::subscribe(&broadcast)
.await
.expect("subscribe to the catalog");

// The relay is serving the catalog before stdin ends, so the finish is queued on a
// live subscription when the process exits.
tokio::time::timeout(TIMEOUT, catalogs.next())
running(&mut child, "the first catalog", catalogs.next())
.await
.expect("catalog timed out")
.expect("catalog read")
.expect("a catalog");

Expand Down
Loading