feat(nests): T16 Phase 2 — real hang-listen + hang-publish bodies

Replaces the Phase 1 stubs with working moq-lite-03 audio
publish/subscribe loops:

- hang-publish: opens a broadcast under <relay>/<broadcast>,
  publishes a hang Catalog with one Opus / Container::Legacy
  audio rendition, and pumps Opus-encoded sine-wave frames in
  groups of 5 (matching Amethyst's NestMoqLiteBroadcaster
  default) for --duration seconds.
- hang-listen: connects to the same broadcast, reads the catalog,
  picks the first Opus audio rendition with container.kind=legacy,
  decodes each Opus packet, and writes Float32 little-endian PCM
  to --output-pcm (or stdout with `-`).

Both use moq-native 0.13 with the quinn + aws-lc-rs features and
explicitly install the rustls aws-lc-rs crypto provider in main()
since rustls 0.23 no longer auto-installs.

Verified Rust↔Rust interop end-to-end: 3-second 440 Hz publish →
2.7 s of decoded PCM, RMS 0.35 (matching half-amplitude sine),
880 zero-crossings/sec (matches 440 Hz exactly).

Phase 2.B (JVM Opus encoder/decoder) and 2.C (I1 amethyst-speaker
→ hang-listen) are next.

https://claude.ai/code/session_01ERJPUYfdLPwZ99pr5EcEcV
This commit is contained in:
Claude
2026-05-06 20:58:09 +00:00
parent 284a203070
commit 33a9e8f053
5 changed files with 3470 additions and 32 deletions
+2967 -6
View File
File diff suppressed because it is too large Load Diff
+13 -5
View File
@@ -5,11 +5,10 @@ edition.workspace = true
publish.workspace = true
license.workspace = true
# Phase 1: stub. Subscribes are added in Phase 2 via `hang` +
# `moq-lite` + `web-transport-quinn` against the rev pinned in
# ../REV. See nestsClient/plans/2026-05-06-cross-stack-interop-test.md
# (Phase 2 step 8) for the catalog → AudioConfig → Container::Legacy
# decode loop this binary will host.
# Real subscribe/decode body: connects to a moq-lite-03 relay, reads
# the hang catalog, picks the first Opus / Container::Legacy audio
# rendition, and writes Float32 PCM to --output-pcm. Used by the
# cross-stack interop tests in nestsClient/src/jvmTest/.../interop/native/.
[[bin]]
name = "hang-listen"
@@ -19,3 +18,12 @@ path = "src/main.rs"
anyhow.workspace = true
clap.workspace = true
tokio.workspace = true
hang = "0.15"
moq-lite = "0.15"
moq-mux = "0.3"
moq-native = { version = "0.13", default-features = false, features = ["quinn", "aws-lc-rs"] }
opus = "0.3"
rustls = { version = "0.23", default-features = false, features = ["aws-lc-rs"] }
tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
url = "2"
+232 -10
View File
@@ -1,36 +1,258 @@
//! hang-listen — reference moq-lite / hang audio listener for the
//! cross-stack interop harness. **Phase 1 stub** — Phase 2 fills in
//! the real subscribe loop. See
//! hang-listen — reference moq-lite / hang audio listener.
//!
//! Used by the Amethyst cross-stack interop test harness to verify
//! that an Amethyst Kotlin speaker is intelligible to the canonical
//! `kixelated/moq` listener stack. See
//! `nestsClient/plans/2026-05-06-cross-stack-interop-test.md`.
//!
//! Wire path:
//! 1. Connect via `web-transport-quinn` over QUIC.
//! 2. Open a `moq-lite-03` session (via moq_native::ClientConfig).
//! 3. Subscribe to the hang Catalog track at `<broadcast>/catalog.json`.
//! 4. Pick the first audio rendition with `codec="opus"` and
//! `container.kind="legacy"`.
//! 5. Subscribe to that rendition's track via
//! `moq_mux::container::Consumer<Hang::Legacy>`.
//! 6. For each frame: decode Opus → Float32 PCM, write to stdout
//! or `--output-pcm` as raw little-endian f32s.
use anyhow::Result;
use std::io::Write;
use std::path::PathBuf;
use std::time::Duration;
use anyhow::{Context, anyhow};
use clap::Parser;
use hang::catalog::{AudioCodec, Container};
const SAMPLE_RATE_HZ: u32 = 48_000;
/// 120 ms at 48 kHz — Opus's worst-case frame size; pre-allocate
/// once and let `Decoder::decode` write what it actually decoded.
const MAX_PCM_PER_PACKET: usize = (SAMPLE_RATE_HZ as usize) / 1000 * 120;
#[derive(Parser, Debug)]
#[command(
name = "hang-listen",
about = "Reference moq-lite / hang audio listener (Phase 1 stub)"
about = "Reference moq-lite / hang audio listener for cross-stack interop"
)]
struct Args {
/// HTTPS URL of the relay, e.g. `https://127.0.0.1:34721`.
#[arg(long)]
relay_url: String,
/// Optional JWT for the `?jwt=` query string. The Amethyst test
/// harness configures the relay with `--auth-public ""`, in which
/// case this can be omitted.
#[arg(long)]
jwt: Option<String>,
/// Broadcast namespace. The full path is `<broadcast>/<track>`.
#[arg(long)]
broadcast: String,
/// Maximum runtime in seconds.
#[arg(long, default_value_t = 5)]
duration: u64,
/// Output Float32 little-endian PCM here. Use `-` for stdout.
/// If absent, the binary discards PCM (used as a smoke test).
#[arg(long)]
output_pcm: Option<String>,
}
#[tokio::main]
async fn main() -> Result<()> {
async fn main() -> anyhow::Result<()> {
// rustls 0.23 requires an explicit crypto-provider install.
// Mirror moq-relay's main.rs choice (aws-lc-rs).
let _ = rustls::crypto::aws_lc_rs::default_provider().install_default();
// Init logger early so config / handshake errors surface.
let _ = tracing_subscriber::fmt()
.with_env_filter(tracing_subscriber::EnvFilter::from_default_env())
.with_writer(std::io::stderr)
.try_init();
let args = Args::parse();
eprintln!(
"hang-listen Phase-1 stub — relay_url={} broadcast={} duration={}s output_pcm={:?}",
args.relay_url, args.broadcast, args.duration, args.output_pcm
let result = tokio::time::timeout(
Duration::from_secs(args.duration + 5),
run(args),
)
.await
.context("hang-listen wallclock timeout")?;
result
}
async fn run(args: Args) -> anyhow::Result<()> {
let url = build_url(&args.relay_url, &args.broadcast, args.jwt.as_deref())?;
// moq-lite-03 ALPN, IPv4 client bind (sandbox friendly), and TLS
// verification disabled so the test harness's --tls-generate
// cert chain works without a custom truststore.
let cfg = moq_native::ClientConfig::parse_from([
"hang-listen",
"--client-bind",
"127.0.0.1:0",
"--client-version",
"moq-lite-03",
"--tls-disable-verify=true",
]);
let client = cfg.init().context("init moq client")?;
// Set up an Origin so the session can publish incoming
// broadcasts to us, then drive both the session and the
// subscribe loop in parallel.
let origin = moq_lite::Origin::produce();
let consumer = origin.consume();
let session_url = url.clone();
let session = tokio::spawn(async move {
// Use reconnect() with a tight timeout so the test exits
// quickly when the relay drops us. closed() returns when the
// backoff loop finally gives up.
let reconnect = client.with_consume(origin).reconnect(session_url);
if let Err(err) = reconnect.closed().await {
tracing::warn!(%err, "reconnect loop exited");
}
});
let listen_result = listen(consumer, args.output_pcm.as_deref(), args.duration).await;
// The session task will exit on its own when the URL closes; we
// don't need to abort it for a clean shutdown.
drop(session);
listen_result
}
async fn listen(
mut origin: moq_lite::OriginConsumer,
output_pcm: Option<&str>,
duration_sec: u64,
) -> anyhow::Result<()> {
// Open the PCM sink up front so we fail fast on a bad path.
let mut pcm: Box<dyn Write + Send> = match output_pcm {
Some("-") => Box::new(std::io::stdout()),
Some(path) => Box::new(
std::fs::File::create(PathBuf::from(path))
.with_context(|| format!("create output-pcm file '{path}'"))?,
),
None => Box::new(std::io::sink()),
};
// Wait for the broadcast to be announced. The relay forwards
// any matching broadcast across the configured namespace.
let (path, broadcast) = origin
.announced()
.await
.ok_or_else(|| anyhow!("origin closed before any broadcast announced"))?;
let broadcast = broadcast.ok_or_else(|| anyhow!("broadcast unannounced: {path}"))?;
tracing::info!(%path, "broadcast announced");
// Subscribe to the catalog and read the first published version.
let catalog_track = broadcast
.subscribe_track(&hang::Catalog::default_track())
.context("subscribe catalog")?;
let mut catalog = hang::CatalogConsumer::new(catalog_track);
let info = catalog
.next()
.await
.context("read catalog")?
.ok_or_else(|| anyhow!("catalog ended before first publish"))?;
// Pick the first Opus / Container::Legacy audio rendition.
let (track_name, audio_cfg) = info
.audio
.renditions
.iter()
.find(|(_, cfg)| matches!(cfg.codec, AudioCodec::Opus) && cfg.container == Container::Legacy)
.ok_or_else(|| {
anyhow!(
"no audio rendition with codec=opus container.kind=legacy in catalog: {:?}",
info.audio.renditions.keys().collect::<Vec<_>>()
)
})?;
// Audio renditions advertise `numberOfChannels` (1 for mono, 2 for
// stereo). nests speakers send mono; stereo is exercised by I4.
let channels = match audio_cfg.channel_count {
1 => opus::Channels::Mono,
2 => opus::Channels::Stereo,
n => anyhow::bail!("unsupported channel count: {n}"),
};
tracing::info!(
track = %track_name,
sample_rate = audio_cfg.sample_rate,
channels = audio_cfg.channel_count,
"subscribing to audio rendition"
);
eprintln!("Phase 2 will implement the actual subscribe loop. Exiting cleanly.");
let track = moq_lite::Track {
name: track_name.clone(),
priority: 1,
};
let track_consumer = broadcast.subscribe_track(&track).context("subscribe audio")?;
let mut frames = moq_mux::hang::Consumer::new(track_consumer, moq_mux::hang::Legacy)
// Zero latency = aggressive skip. We prefer a more forgiving
// budget so jitter doesn't drop frames in tests.
.with_latency(Duration::from_millis(500));
let mut decoder =
opus::Decoder::new(audio_cfg.sample_rate, channels).context("init opus decoder")?;
let mut pcm_buf = vec![0i16; MAX_PCM_PER_PACKET * audio_cfg.channel_count as usize];
let deadline = tokio::time::Instant::now() + Duration::from_secs(duration_sec);
let mut total_samples: u64 = 0;
let mut frame_count: u64 = 0;
loop {
let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
if remaining.is_zero() {
break;
}
let frame = match tokio::time::timeout(remaining, frames.read()).await {
Ok(Ok(Some(f))) => f,
Ok(Ok(None)) => {
tracing::info!("track ended");
break;
}
Ok(Err(e)) => return Err(anyhow::Error::new(e).context("read audio frame")),
Err(_) => {
tracing::info!("duration elapsed");
break;
}
};
// payload is the raw Opus packet — the timestamp varint has
// already been stripped by `Hang::Legacy` decoding.
let n = decoder
.decode(&frame.payload, &mut pcm_buf, false)
.with_context(|| format!("decode opus packet ({} bytes)", frame.payload.len()))?;
// n is the number of samples per channel; total interleaved
// samples in pcm_buf is n * channels.
let interleaved = n * audio_cfg.channel_count as usize;
for s in &pcm_buf[..interleaved] {
// i16 → f32 in [-1, 1].
let f = (*s as f32) / 32_768.0;
pcm.write_all(&f.to_le_bytes())
.context("write pcm sample")?;
}
total_samples += interleaved as u64;
frame_count += 1;
}
pcm.flush().ok();
tracing::info!(frames = frame_count, samples = total_samples, "hang-listen done");
Ok(())
}
fn build_url(relay_url: &str, broadcast: &str, jwt: Option<&str>) -> anyhow::Result<url::Url> {
let trimmed = relay_url.trim_end_matches('/');
let raw = if let Some(jwt) = jwt {
format!("{trimmed}/{broadcast}?jwt={jwt}")
} else {
format!("{trimmed}/{broadcast}")
};
url::Url::parse(&raw).with_context(|| format!("malformed relay/broadcast url: {raw}"))
}
+14 -2
View File
@@ -5,8 +5,10 @@ edition.workspace = true
publish.workspace = true
license.workspace = true
# Phase 1: stub. Phase 2 wires this against `hang` + `audiopus` for
# the reverse direction (Rust → Amethyst) interop scenarios.
# Reference moq-lite / hang publisher: opens a broadcast, declares one
# Opus / Container::Legacy audio rendition, encodes a sine wave, and
# pumps frames in 5-frame groups for `--duration` seconds. Used by
# the cross-stack interop tests for the Rust → Amethyst direction.
[[bin]]
name = "hang-publish"
@@ -16,3 +18,13 @@ path = "src/main.rs"
anyhow.workspace = true
clap.workspace = true
tokio.workspace = true
hang = "0.15"
moq-lite = "0.15"
moq-native = { version = "0.13", default-features = false, features = ["quinn", "aws-lc-rs"] }
opus = "0.3"
bytes = "1"
rustls = { version = "0.23", default-features = false, features = ["aws-lc-rs"] }
serde_json = "1"
tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
url = "2"
+244 -9
View File
@@ -1,36 +1,271 @@
//! hang-publish — reference moq-lite / hang audio publisher for the
//! cross-stack interop harness. **Phase 1 stub.**
//! hang-publish — reference moq-lite / hang audio publisher.
//!
//! Opens a broadcast at `<relay>/<broadcast>`, publishes a hang
//! catalog with one Opus / Container::Legacy audio rendition, and
//! pumps Opus-encoded sine-wave frames in groups of 5 for
//! `--duration` seconds. Used by the cross-stack interop tests for
//! the Rust → Amethyst direction. See
//! `nestsClient/plans/2026-05-06-cross-stack-interop-test.md`.
use anyhow::Result;
use std::time::Duration;
use anyhow::{Context, anyhow};
use bytes::Bytes;
use clap::Parser;
use hang::catalog::{Audio, AudioCodec, AudioConfig, Catalog, Container};
const SAMPLE_RATE_HZ: u32 = 48_000;
/// 20 ms at 48 kHz — same frame size Amethyst speakers send.
const FRAME_SIZE_SAMPLES: usize = 960;
/// Microseconds per frame: 20_000 = 1_000_000 * 960 / 48_000.
const FRAME_DURATION_US: u64 = 20_000;
/// 5 frames per group → 100 ms group cadence, matching nests speaker.
const FRAMES_PER_GROUP: usize = 5;
/// Audio rendition track name in the catalog.
const TRACK_NAME: &str = "audio";
#[derive(Parser, Debug)]
#[command(
name = "hang-publish",
about = "Reference moq-lite / hang audio publisher (Phase 1 stub)"
about = "Reference moq-lite / hang audio publisher for cross-stack interop"
)]
struct Args {
/// HTTPS URL of the relay, e.g. `https://127.0.0.1:34721`.
#[arg(long)]
relay_url: String,
/// Optional JWT for the `?jwt=` query string.
#[arg(long)]
jwt: Option<String>,
/// Broadcast namespace (path under the relay root).
#[arg(long)]
broadcast: String,
/// Sine-wave frequency in Hz (mono only). For stereo, see
/// `--freq-hz-right` once that lands.
#[arg(long, default_value_t = 440)]
freq_hz: u32,
/// Maximum runtime in seconds.
#[arg(long, default_value_t = 5)]
duration: u64,
/// Channel count: 1 (mono) or 2 (stereo). Stereo uses the same
/// frequency on both channels for now; per-channel frequency is
/// a Phase-2 follow-up.
#[arg(long, default_value_t = 1)]
channels: u32,
}
#[tokio::main]
async fn main() -> Result<()> {
async fn main() -> anyhow::Result<()> {
// rustls 0.23 requires an explicit crypto-provider install.
let _ = rustls::crypto::aws_lc_rs::default_provider().install_default();
let _ = tracing_subscriber::fmt()
.with_env_filter(tracing_subscriber::EnvFilter::from_default_env())
.with_writer(std::io::stderr)
.try_init();
let args = Args::parse();
eprintln!(
"hang-publish Phase-1 stub — relay_url={} broadcast={} freq_hz={} duration={}s channels={}",
args.relay_url, args.broadcast, args.freq_hz, args.duration, args.channels
let result = tokio::time::timeout(
Duration::from_secs(args.duration + 5),
run(args),
)
.await
.context("hang-publish wallclock timeout")?;
result
}
async fn run(args: Args) -> anyhow::Result<()> {
let url = build_url(&args.relay_url, &args.broadcast, args.jwt.as_deref())?;
let cfg = moq_native::ClientConfig::parse_from([
"hang-publish",
"--client-bind",
"127.0.0.1:0",
"--client-version",
"moq-lite-03",
"--tls-disable-verify=true",
]);
let client = cfg.init().context("init moq client")?;
// Set up the publish side: a producer this binary writes to,
// and a consumer the moq-native client forwards to the relay.
let origin = moq_lite::Origin::produce();
let publish_consumer = origin.consume();
// Start the reconnect loop in the background. It owns the
// session lifecycle.
let session_url = url.clone();
let session = tokio::spawn(async move {
let reconnect = client.with_publish(publish_consumer).reconnect(session_url);
if let Err(err) = reconnect.closed().await {
tracing::warn!(%err, "reconnect loop exited");
}
});
// Result of the publish loop is what determines test pass/fail.
let publish_result = publish(&origin, &args).await;
// Once we drop origin all published broadcasts unannounce; the
// reconnect task exits when the session closes.
drop(session);
publish_result
}
async fn publish(origin: &moq_lite::OriginProducer, args: &Args) -> anyhow::Result<()> {
let mut broadcast = origin
.create_broadcast(args.broadcast.as_str())
.ok_or_else(|| anyhow!("broadcast '{}' not allowed by origin", args.broadcast))?;
// 1. Catalog track. We declare one Opus rendition; the JSON
// payload mirrors what Amethyst's MoqLiteHangCatalog produces.
let mut catalog_track = broadcast
.create_track(hang::Catalog::default_track())
.context("create catalog track")?;
let mut renditions = std::collections::BTreeMap::new();
renditions.insert(
TRACK_NAME.to_string(),
AudioConfig {
codec: AudioCodec::Opus,
sample_rate: SAMPLE_RATE_HZ,
channel_count: args.channels,
bitrate: Some(32_000),
description: None,
container: Container::Legacy,
jitter: None,
},
);
let catalog = Catalog {
audio: Audio { renditions },
..Default::default()
};
let catalog_json =
serde_json::to_vec(&catalog).context("serialize catalog json")?;
let mut catalog_group = catalog_track
.create_group(moq_lite::Group { sequence: 0 })
.context("create catalog group")?;
catalog_group
.write_frame(catalog_json)
.context("publish catalog frame")?;
catalog_group.finish().ok();
// We don't finish() the catalog_track itself yet — moq-lite
// treats track-end as broadcast-end, and we want the audio
// track to keep streaming.
// 2. Audio track.
let mut audio_track = broadcast
.create_track(moq_lite::Track {
name: TRACK_NAME.to_string(),
priority: 1,
})
.context("create audio track")?;
let channels = match args.channels {
1 => opus::Channels::Mono,
2 => opus::Channels::Stereo,
n => anyhow::bail!("unsupported channel count: {n}"),
};
let mut encoder = opus::Encoder::new(SAMPLE_RATE_HZ, channels, opus::Application::Audio)
.context("init opus encoder")?;
encoder
.set_bitrate(opus::Bitrate::Bits(32_000))
.context("set opus bitrate")?;
let total_frames = (args.duration * 1_000_000 / FRAME_DURATION_US) as usize;
let phase_step =
2.0_f64 * std::f64::consts::PI * (args.freq_hz as f64) / (SAMPLE_RATE_HZ as f64);
let mut sample_idx: u64 = 0;
// Sized to libopus's worst-case output for one 20 ms frame.
let mut opus_buf = vec![0u8; 4_000];
let frame_period = Duration::from_micros(FRAME_DURATION_US);
let mut next_send = tokio::time::Instant::now();
let mut group_idx: u64 = 0;
let mut frames_in_group = 0usize;
let mut group: Option<moq_lite::GroupProducer> = None;
for frame_no in 0..total_frames {
// Generate one PCM frame at the configured frequency.
let mut pcm = vec![0i16; FRAME_SIZE_SAMPLES * args.channels as usize];
for i in 0..FRAME_SIZE_SAMPLES {
let t = sample_idx + i as u64;
let v = ((t as f64) * phase_step).sin();
let s = (v * 16_383.0) as i16;
for ch in 0..(args.channels as usize) {
pcm[i * args.channels as usize + ch] = s;
}
}
sample_idx += FRAME_SIZE_SAMPLES as u64;
let n = encoder
.encode(&pcm, &mut opus_buf)
.context("encode opus packet")?;
let opus_packet = Bytes::copy_from_slice(&opus_buf[..n]);
// Wrap the Opus packet in a hang Legacy frame: VarInt
// timestamp prefix + raw codec payload.
let frame = hang::container::Frame {
timestamp: hang::container::Timestamp::from_micros(
(frame_no as u64) * FRAME_DURATION_US,
)
.context("frame timestamp out of range")?,
payload: opus_packet.into(),
};
// Start a new group every FRAMES_PER_GROUP frames. The 5-frame
// group cadence matches Amethyst's NestMoqLiteBroadcaster default
// and produces ~100 ms groups.
if frames_in_group == 0 {
if let Some(mut g) = group.take() {
g.finish().ok();
}
group = Some(
audio_track
.create_group(moq_lite::Group { sequence: group_idx })
.context("create audio group")?,
);
group_idx += 1;
}
let g = group.as_mut().expect("group always Some after init");
frame.encode(g).context("encode hang frame into group")?;
frames_in_group += 1;
if frames_in_group == FRAMES_PER_GROUP {
frames_in_group = 0;
}
next_send += frame_period;
tokio::time::sleep_until(next_send).await;
}
if let Some(mut g) = group.take() {
g.finish().ok();
}
audio_track.finish().ok();
catalog_track.finish().ok();
tracing::info!(
frames = total_frames,
groups = group_idx,
"hang-publish done"
);
eprintln!("Phase 2 will implement the actual publish loop. Exiting cleanly.");
Ok(())
}
fn build_url(relay_url: &str, broadcast: &str, jwt: Option<&str>) -> anyhow::Result<url::Url> {
let trimmed = relay_url.trim_end_matches('/');
let raw = if let Some(jwt) = jwt {
format!("{trimmed}/{broadcast}?jwt={jwt}")
} else {
format!("{trimmed}/{broadcast}")
};
url::Url::parse(&raw).with_context(|| format!("malformed relay/broadcast url: {raw}"))
}