From 17abeed7aed697a66589f3863a26fc69e43e37d6 Mon Sep 17 00:00:00 2001 From: Breadway Date: Sun, 16 Aug 2026 14:15:56 +0800 Subject: [PATCH 1/2] Fix Cast teardown leaks, keyframe latch, and DLNA session lifecycle Dropped frames never requested a keyframe, SessionEnded skipped ordered stop (portal/TV/FFI leak, next start could abort), and a late end could kill the following cast. Failed starts left PlatformClientPosix alive. DLNA leaked its HTTP server and ignored portal EOS. Also: start no longer blocks the daemon actor, IPC accept/request loops stay up, HLS Range is clamped, LAN IP follows the renderer subnet, and the picker closes before the portal dialog and handles Escape. --- .forgejo/workflows/check.yml | 2 +- .forgejo/workflows/dev-release.yml | 11 +- AGENTS.md | 18 +-- EVENTS.md | 2 +- bakery.toml | 5 +- breadcast-caststream-sys/src/facade.cc | 77 ++++++++-- breadcast-caststream-sys/src/facade.h | 21 +-- breadcast-caststream-sys/src/lib.rs | 4 +- breadcast-core/src/capture/portal.rs | 4 + breadcast-core/src/caststream.rs | 22 +-- breadcast-core/src/discovery.rs | 20 ++- breadcast-core/src/dlna/discovery.rs | 8 +- breadcast-core/src/dlna/mod.rs | 4 + breadcast-core/src/dlna/session.rs | 37 ++++- breadcast-core/src/http_server.rs | 62 +++++++- breadcast-core/src/net.rs | 85 +++++++++-- breadcast-core/src/pipeline/mod.rs | 66 ++++++-- breadcast/src/css.rs | 5 + breadcast/src/ipc_client.rs | 2 + breadcast/src/main.rs | 75 ++++++++- breadcastd/src/cast_mirror.rs | 118 ++++++++++---- breadcastd/src/daemon.rs | 181 ++++++++++++++++------ breadcastd/src/dlna_mirror.rs | 204 ++++++++++++++++++++----- breadcastd/src/ipc.rs | 35 ++++- contrib/breadcastd.service | 2 +- vendor/rust_cast-0.21.0/PATCHES.md | 13 +- 26 files changed, 861 insertions(+), 222 deletions(-) diff --git a/.forgejo/workflows/check.yml b/.forgejo/workflows/check.yml index b547c34..65ae3e6 100644 --- a/.forgejo/workflows/check.yml +++ b/.forgejo/workflows/check.yml @@ -4,7 +4,7 @@ name: check # main and triggers a dev-track release build. on: push: - branches: ['feature/**', 'fix/**'] + branches: ['main', 'feature/**', 'fix/**'] jobs: check: diff --git a/.forgejo/workflows/dev-release.yml b/.forgejo/workflows/dev-release.yml index 3d36231..0207d3d 100644 --- a/.forgejo/workflows/dev-release.yml +++ b/.forgejo/workflows/dev-release.yml @@ -1,9 +1,8 @@ name: dev release -# Publishes a dev-track build on every push to `dev` — -# separate from release.yml's tag-triggered stable releases. See -# bread-ecosystem's docs/release-channels.md for the three-track policy -# this is part of. +# Publishes a dev-track build on every push to `main` (single-trunk). +# The leftover `dev` trigger is kept so an old branch push still works +# until that branch is deleted. See CONTRIBUTING.md. on: push: branches: ['main', 'dev'] @@ -16,7 +15,7 @@ jobs: run: | set -euo pipefail rm -rf src && mkdir src - git clone --branch main --depth 1 \ + git clone --branch "${GITHUB_REF_NAME}" --depth 1 \ "https://git.breadway.dev/${GITHUB_REPOSITORY}.git" src - name: build @@ -37,7 +36,7 @@ jobs: if [ -n "${LATEST_TAG}" ]; then CUR="${LATEST_TAG}" else - CUR="$(grep -m1 '^version' breadcast/Cargo.toml | sed -E 's/.*"(.*)".*/\1/')" + CUR="$(grep -m1 '^version' Cargo.toml | sed -E 's/.*"(.*)".*/\1/')" fi IFS='.' read -r MA MI PA <<< "${CUR}" SHA="$(git rev-parse --short HEAD)" diff --git a/AGENTS.md b/AGENTS.md index 29ba49e..0a3ec7e 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -22,13 +22,11 @@ out of sync with `dev`/`beta` across most repos in this ecosystem. ## CI - `release.yml` triggers on `push: tags: ['v*']` — tag a release to cut the signed stable build. -- Leftover: `dev-release.yml` still triggers on `push: branches: ['dev']` - and `beta-release.yml` still triggers on `push: branches: ['beta']`. - Those files are leftover from the three-branch model. Do **not** rewrite - them unless you are deliberately migrating CI; document and follow - single-trunk `main` + RC tags instead. -- `check.yml` runs clippy + test on `feature/**` and `fix/**`. -- No separate lint/PR-check pipeline on ordinary commits to `main`. +- `dev-release.yml` publishes a dev-track build on push to `main` (and + still on leftover `dev` so an old branch push does not go silent). +- `beta-release.yml` still triggers on leftover `push: branches: ['beta']`. + Do **not** use that branch; cut beta from a `vX.Y.Z-rc.N` tag instead. +- `check.yml` runs clippy + test on `main`, `feature/**`, and `fix/**`. ## Product cut breadcast is a **bakery product**, not shipped on the BOS ISO, and not part @@ -41,9 +39,9 @@ The daemon + GTK picker + Cast Streaming / DLNA pipelines are built and validated against real hardware. `EVENTS.md` is the bus contract. ## Pins -`bread-theme` and `bread-utils` pin -`git.breadway.dev/Breadway/bread-ecosystem` tag `v0.7.1`. Do not switch -those to github.com or `branch = "main"`. +`bread-theme` pins `git.breadway.dev/Breadway/bread-ecosystem` tag +`v0.7.4` (per-monitor palette). `bread-utils` pins the same repo at +`v0.7.2`. Do not switch those to github.com or `branch = "main"`. ## Don't - Don't embed credentials in remote URLs — SSH or a credential helper only. diff --git a/EVENTS.md b/EVENTS.md index d198d1a..6b76d57 100644 --- a/EVENTS.md +++ b/EVENTS.md @@ -16,7 +16,7 @@ discovery-driven emit half. | Event | Data | When | |-------|------|------| -| `bread.cast.device_found` | `{ "id": "", "name": "", "model": "", "protocol": "cast" \| "dlna" }` | A Chromecast/Google TV (mDNS) or DLNA/UPnP media renderer (SSDP) is discovered (or re-resolved) on the LAN. Fires on every re-resolution, not just the first sighting — treat it as an upsert keyed by `id`, not an append-only log. `id` is protocol-specific and only unique *within* a protocol — a Cast device's mDNS id and a DLNA device's description URL share no namespace. | +| `bread.cast.device_found` | `{ "id": "", "name": "", "model": "", "protocol": "cast" \| "dlna" }` | A Chromecast/Google TV (mDNS) or DLNA/UPnP media renderer (SSDP) is discovered, or an already-known device changes (new host, new name). Re-resolutions that carry the same payload are **not** re-emitted — treat it as an upsert keyed by `id`, not an append-only log. `id` is protocol-specific and only unique *within* a protocol — a Cast device's mDNS id and a DLNA device's description URL share no namespace. | | `bread.cast.mirroring_started` | `{ "device_id": "", "device_name": "", "protocol": "cast" \| "dlna" }` | A mirroring session successfully started, whether triggered by `bread.command.cast.start`, the `breadcast` GTK popup, or (once wired) any other IPC client. | | `bread.cast.mirroring_stopped` | `{}` | A mirroring session ended, whether via an explicit stop (command, IPC, or GTK popup) or unprompted (the portal picker's "stop sharing", the receiver dropping the connection, a DLNA renderer stopping playback from its own remote). There is no separate "stopped by whom" distinction in this event — `breadcastd`'s own logs have that detail if needed. | | `bread.cast.mirroring_failed` | `{ "device_id": "", "error": "" }` | A `start_cast`/`bread.command.cast.start` attempt failed before a session was established (device unreachable, portal capture denied, negotiation timeout, renderer rejected the stream, etc). | diff --git a/bakery.toml b/bakery.toml index 68e5c5d..a6416b0 100644 --- a/bakery.toml +++ b/bakery.toml @@ -1,7 +1,9 @@ name = "breadcast" description = "Cast your screen to any Chromecast/Google TV or DLNA renderer — daemon + GTK4 popup" binaries = ["breadcast", "breadcastd"] -# gst-plugin-va: the `vah264enc` element both encode pipelines use (see +# gst-plugins-base: videoconvert/videoscale/videorate/appsink -- the +# `gstreamer` package is only the core library. gst-plugin-va: the +# `vah264enc` element both encode pipelines use (see # breadcast-core/src/pipeline/mod.rs) -- a separate Arch package from # gst-plugins-bad itself, not bundled into it. jsoncpp/openssl: runtime # shared-library deps of breadcastd itself (not just a build dep of @@ -14,6 +16,7 @@ system_deps = [ "gtk4", "gtk4-layer-shell", "gstreamer", + "gst-plugins-base", "gst-plugin-pipewire", "gst-plugins-bad", "gst-plugin-hlssink3", diff --git a/breadcast-caststream-sys/src/facade.cc b/breadcast-caststream-sys/src/facade.cc index 3ec4288..4176a31 100644 --- a/breadcast-caststream-sys/src/facade.cc +++ b/breadcast-caststream-sys/src/facade.cc @@ -209,9 +209,20 @@ CastStreamSender* breadcast_caststream_sender_create( // than once process-wide only if ShutDown() was called first -- breadcast // only ever has one active cast-streaming session at a time, so this // assumption (baked into PlatformClientPosix's own singleton design) holds. + // + // Create() itself OSP_CHECKs that no instance exists and aborts the + // process if a previous sender leaked (a failed start that never reached + // destroy()). Fail the new create instead of taking the daemon down. + if (openscreen::PlatformClientPosix::GetInstance() != nullptr) { + return nullptr; + } openscreen::PlatformClientPosix::Create(std::chrono::milliseconds(50)); - openscreen::TaskRunner& task_runner = - openscreen::PlatformClientPosix::GetInstance()->GetTaskRunner(); + openscreen::PlatformClientPosix* instance = + openscreen::PlatformClientPosix::GetInstance(); + if (instance == nullptr) { + return nullptr; + } + openscreen::TaskRunner& task_runner = instance->GetTaskRunner(); breadcast_caststream::VideoParams params; params.width = width; @@ -255,6 +266,13 @@ CastStreamSender* breadcast_caststream_sender_create( } void breadcast_caststream_sender_negotiate(CastStreamSender* sender) { + // environment is reset during destroy() on the TaskRunner thread; a call + // that races teardown (or arrives after a failed start) must not + // dereference the unique_ptr. + if (!sender || sender->shutting_down.load(std::memory_order_acquire) || + !sender->environment) { + return; + } sender->environment->task_runner().PostTask([sender] { if (sender->shutting_down.load(std::memory_order_acquire) || !sender->session) { @@ -272,6 +290,10 @@ void breadcast_caststream_sender_on_message(CastStreamSender* sender, size_t message_namespace_len, const char* message, size_t message_len) { + if (!sender || sender->shutting_down.load(std::memory_order_acquire) || + !sender->environment) { + return; + } auto source = std::make_shared(source_id, source_id_len); auto ns = std::make_shared(message_namespace, message_namespace_len); auto body = std::make_shared(message, message_len); @@ -289,6 +311,10 @@ int32_t breadcast_caststream_sender_enqueue_frame(CastStreamSender* sender, size_t data_len, int32_t is_key_frame, int64_t capture_time_us) { + if (!sender || sender->shutting_down.load(std::memory_order_acquire) || + !sender->environment) { + return -1; + } if (!sender->negotiated.load(std::memory_order_acquire)) { return -1; } @@ -378,6 +404,12 @@ int32_t breadcast_caststream_sender_enqueue_frame(CastStreamSender* sender, switch (video_sender->EnqueueFrame(frame)) { case Sender::OK: sender->enqueue_ok.fetch_add(1, std::memory_order_relaxed); + // A key frame that actually landed resyncs the decoder, so the + // "next frame must be an IDR" latch can clear. A rejected key + // frame leaves the flag set (the reject branches below). + if (is_key) { + sender->frame_chain_broken.store(false, std::memory_order_relaxed); + } break; case Sender::PAYLOAD_TOO_LARGE: sender->enqueue_payload_too_large.fetch_add(1, std::memory_order_relaxed); @@ -418,7 +450,18 @@ void breadcast_caststream_sender_take_stats(CastStreamSender* sender, } int32_t breadcast_caststream_sender_needs_key_frame(CastStreamSender* sender) { - return sender->needs_key_frame.load(std::memory_order_relaxed) ? 1 : 0; + if (!sender) { + return 0; + } + // `frame_chain_broken` is the drop-side latch SchedulePoll must not + // clobber (see the field comment). Load, don't exchange: the encoder + // may take several frames to honour a force-key-unit, and a consuming + // read here would forget the drop if the next poll happened before + // that IDR was actually EnqueueFrame'd (cleared on OK + is_key above). + return (sender->frame_chain_broken.load(std::memory_order_relaxed) || + sender->needs_key_frame.load(std::memory_order_relaxed)) + ? 1 + : 0; } int32_t breadcast_caststream_sender_estimated_bandwidth_bps(CastStreamSender* sender) { @@ -437,18 +480,24 @@ void breadcast_caststream_sender_destroy(CastStreamSender* sender) { // These must be torn down on the TaskRunner thread (they hold raw // references into it and into `environment`), so hop over there and block - // until it's done before shutting the TaskRunner itself down. - std::promise done; - std::future done_future = done.get_future(); - sender->environment->task_runner().PostTask([sender, &done] { - sender->session.reset(); - sender->message_port.reset(); - sender->environment.reset(); - done.set_value(); - }); - done_future.wait(); + // until it's done before shutting the TaskRunner itself down. A sender + // whose Environment never got built (create failed mid-flight) has no + // task runner to hop to. + if (sender->environment) { + std::promise done; + std::future done_future = done.get_future(); + sender->environment->task_runner().PostTask([sender, &done] { + sender->session.reset(); + sender->message_port.reset(); + sender->environment.reset(); + done.set_value(); + }); + done_future.wait(); + } - openscreen::PlatformClientPosix::ShutDown(); + if (openscreen::PlatformClientPosix::GetInstance() != nullptr) { + openscreen::PlatformClientPosix::ShutDown(); + } delete sender; } diff --git a/breadcast-caststream-sys/src/facade.h b/breadcast-caststream-sys/src/facade.h index 8411721..840d890 100644 --- a/breadcast-caststream-sys/src/facade.h +++ b/breadcast-caststream-sys/src/facade.h @@ -88,15 +88,18 @@ void breadcast_caststream_sender_on_message(CastStreamSender* sender, const char* message, size_t message_len); -// Enqueues one encoded video access unit (Annex-B H.264) for sending. -// `data` is copied before this returns, so the caller may reuse/free its -// buffer immediately after. `capture_time_us` is only used to derive the -// RTP timestamp's relative spacing between frames (it does not need to be -// wall-clock-accurate, just monotonically increasing and proportional to -// real elapsed time between frames). Returns 0 if queued, nonzero if the -// session isn't negotiated yet or the frame was rejected (e.g. too large, -// or the in-flight queue is full -- the caller should back off encoding -// when this happens rather than treating it as fatal). +// Posts one encoded video access unit (Annex-B H.264) onto openscreen's +// TaskRunner for sending. `data` is copied before this returns, so the +// caller may reuse/free its buffer immediately after. `capture_time_us` is +// only used to derive the RTP timestamp's relative spacing between frames +// (it does not need to be wall-clock-accurate, just monotonically +// increasing and proportional to real elapsed time between frames). +// +// Returns 0 if the frame was *posted* (not necessarily accepted -- +// Sender::EnqueueFrame runs later on the TaskRunner and may still reject +// it), or nonzero if the session isn't negotiated yet / is shutting down. +// Accept/reject outcomes are visible only via take_stats() and, for +// dropped frames that break the H.264 reference chain, needs_key_frame(). int32_t breadcast_caststream_sender_enqueue_frame(CastStreamSender* sender, const uint8_t* data, size_t data_len, diff --git a/breadcast-caststream-sys/src/lib.rs b/breadcast-caststream-sys/src/lib.rs index 5837002..222bb6f 100644 --- a/breadcast-caststream-sys/src/lib.rs +++ b/breadcast-caststream-sys/src/lib.rs @@ -1,8 +1,8 @@ //! Raw FFI bindings to `src/facade.h`/`src/facade.cc`, which wrap a pruned, //! vendored subset of `chromium/openscreen`'s Cast Streaming sender (see //! `vendor/openscreen/PATCHES.md`). This crate is intentionally low-level and -//! unsafe -- see `breadcast-caststream` (not this crate) for the ergonomic, -//! thread-safe wrapper most callers should use instead. +//! unsafe -- see `breadcast-core::caststream` for the ergonomic, thread-safe +//! wrapper most callers should use instead. //! //! # Threading contract //! diff --git a/breadcast-core/src/capture/portal.rs b/breadcast-core/src/capture/portal.rs index 77b99cc..4241c97 100644 --- a/breadcast-core/src/capture/portal.rs +++ b/breadcast-core/src/capture/portal.rs @@ -117,6 +117,10 @@ impl CaptureSession { .open(&token_path) { use std::io::Write; + use std::os::unix::fs::PermissionsExt; + // `mode()` only applies on create. A pre-existing + // world-readable token would keep its mode otherwise. + let _ = file.set_permissions(std::fs::Permissions::from_mode(0o600)); let _ = file.write_all(token.as_bytes()); } } diff --git a/breadcast-core/src/caststream.rs b/breadcast-core/src/caststream.rs index 433efa1..9815cfd 100644 --- a/breadcast-core/src/caststream.rs +++ b/breadcast-core/src/caststream.rs @@ -225,16 +225,18 @@ impl CastStreamSender { } } - /// Enqueues one encoded video access unit (Annex-B H.264) for sending. - /// `capture_time_us` only needs to be monotonically increasing and - /// proportional to real elapsed time between frames -- it does not need - /// to be wall-clock-accurate. + /// Posts one encoded video access unit (Annex-B H.264) onto the C++ + /// TaskRunner. `capture_time_us` only needs to be monotonically + /// increasing and proportional to real elapsed time between frames -- + /// it does not need to be wall-clock-accurate. /// - /// Returns an error if the session isn't negotiated yet or the frame - /// was rejected under backpressure; callers should treat the latter as - /// a dropped frame, not a fatal condition (see - /// [`Self::needs_key_frame`]/[`Self::estimated_bandwidth_bps`] for how - /// to react). + /// Returns an error only if the session isn't negotiated (or is + /// shutting down). A `Ok(())` means the frame was *posted*, not that + /// `Sender::EnqueueFrame` accepted it -- accept/reject is visible via + /// [`Self::enqueue_stats`], and a reject that breaks the H.264 + /// reference chain latches [`Self::needs_key_frame`]. Treating this + /// return as accept/reject is how an earlier version hid a frozen + /// picture behind a healthy "30fps enqueued" log line. pub fn enqueue_frame(&self, data: &[u8], is_key_frame: bool, capture_time_us: i64) -> Result<()> { let result = unsafe { breadcast_caststream_sender_enqueue_frame( @@ -246,7 +248,7 @@ impl CastStreamSender { ) }; if result != 0 { - bail!("frame not enqueued (session not negotiated yet)"); + bail!("frame not posted (session not negotiated or shutting down)"); } Ok(()) } diff --git a/breadcast-core/src/discovery.rs b/breadcast-core/src/discovery.rs index 35b8cda..00a8495 100644 --- a/breadcast-core/src/discovery.rs +++ b/breadcast-core/src/discovery.rs @@ -29,7 +29,7 @@ const ADDRESS_STALE_AFTER: Duration = Duration::from_secs(300); /// to two IPv4 addresses (e.g. wired + wireless, or a DHCP lease change), /// an unordered choice can silently pick a stale/unreachable one, and pick /// a *different* one across otherwise-identical runs. -fn pick_address(addresses: &[(ScopedIp, Instant)]) -> Option { +fn pick_address(addresses: &[(ScopedIp, Instant)]) -> Option { let now = Instant::now(); let mut candidates: Vec<&(ScopedIp, Instant)> = addresses .iter() @@ -50,7 +50,10 @@ fn pick_address(addresses: &[(ScopedIp, Instant)]) -> Option { }) }) .or_else(|| candidates.first()) - .map(|(ip, _)| ip.to_ip_addr()) + // Prefer ScopedIp's Display so a last-resort link-local IPv6 + // keeps its `%iface` zone -- `IpAddr` drops it and + // `TcpStream::connect("fe80::…")` then fails with EINVAL. + .map(|(ip, _)| ip.to_string()) } const SERVICE_TYPE: &str = "_googlecast._tcp.local."; @@ -140,10 +143,21 @@ impl Discovery { } } - let Some(host) = pick_address(addrs).map(|ip| ip.to_string()) else { + let Some(host) = pick_address(addrs) else { warn!(%id, "cast device resolved with no usable addresses, skipping"); continue; }; + // First ServiceResolved is often IPv6-only (or a + // zoneless fe80::). rust_cast / Cast Streaming + // cannot connect to that. Wait for a later + // resolve that carries IPv4 or a scoped address. + if host.parse::().is_ok() + && host.starts_with("fe80:") + && !host.contains('%') + { + warn!(%id, %host, "cast device resolved to an unscoped link-local IPv6, waiting for a better address"); + continue; + } fullname_to_id.insert(info.get_fullname().to_string(), id.clone()); diff --git a/breadcast-core/src/dlna/discovery.rs b/breadcast-core/src/dlna/discovery.rs index 9e44a07..79c1fef 100644 --- a/breadcast-core/src/dlna/discovery.rs +++ b/breadcast-core/src/dlna/discovery.rs @@ -6,7 +6,7 @@ use tokio::sync::mpsc; use tracing::{debug, warn}; use super::device::DlnaDevice; -use super::AV_TRANSPORT; +use super::{AV_TRANSPORT, MEDIA_RENDERER}; /// How often to re-issue an SSDP search burst. Unlike mDNS (continuous /// multicast browsing via `mdns-sd` in [`crate::discovery`]), SSDP has no @@ -69,7 +69,7 @@ impl DlnaDiscovery { let mut known: HashMap = HashMap::new(); loop { - let search_target = SearchTarget::URN(AV_TRANSPORT); + let search_target = SearchTarget::URN(MEDIA_RENDERER); match rupnp::discover(&search_target, SEARCH_TIMEOUT, None).await { Ok(stream) => { use futures_util::StreamExt; @@ -86,6 +86,10 @@ impl DlnaDiscovery { }; let url = device.url().to_string(); + if device.find_service(&AV_TRANSPORT).is_none() { + debug!(%url, "UPnP device has no AVTransport, skipping"); + continue; + } confirmed.insert(url.clone()); if let Some(entry) = known.get_mut(&url) { diff --git a/breadcast-core/src/dlna/mod.rs b/breadcast-core/src/dlna/mod.rs index 26b2199..e6d78f8 100644 --- a/breadcast-core/src/dlna/mod.rs +++ b/breadcast-core/src/dlna/mod.rs @@ -21,3 +21,7 @@ pub use discovery::{DlnaDiscovery, DlnaDiscoveryEvent}; pub use session::DlnaSession; const AV_TRANSPORT: URN = URN::service("schemas-upnp-org", "AVTransport", 1); +/// Device-type search. Many TVs answer M-SEARCH for MediaRenderer but not +/// for the AVTransport *service* URN (Windows "Cast to Device" searches +/// this). After resolve we still require AVTransport. +const MEDIA_RENDERER: URN = URN::device("schemas-upnp-org", "MediaRenderer", 1); diff --git a/breadcast-core/src/dlna/session.rs b/breadcast-core/src/dlna/session.rs index 2447ce3..31b0e1d 100644 --- a/breadcast-core/src/dlna/session.rs +++ b/breadcast-core/src/dlna/session.rs @@ -56,18 +56,36 @@ impl DlnaSession { /// renderer would reject with no useful diagnostic. pub async fn load(&self, content_url: &str) -> Result<()> { let escaped = xml_escape(content_url); + // Many Samsung/LG renderers reject an empty CurrentURIMetaData and + // need DIDL-Lite + protocolInfo before they will play a live HLS + // playlist. The DIDL itself is then XML-escaped for the SOAP body. + let didl = format!( + r#"breadcastobject.item.videoItem{escaped}"# + ); + let metadata = xml_escape(&didl); let set_uri_payload = format!( - "0{escaped}" + "0{escaped}{metadata}" ); self.service .action(&self.device_url, "SetAVTransportURI", &set_uri_payload) .await .context("SetAVTransportURI failed")?; - self.service + // Some renderers auto-play after SetURI and then reject Play; + // others stay TRANSITIONING for a moment. Either PLAYING state is + // success. + match self + .service .action(&self.device_url, "Play", "01") .await - .context("Play failed")?; + { + Ok(_) => {} + Err(e) => { + if self.transport_state().await.ok().as_deref() != Some("PLAYING") { + return Err(e).context("Play failed"); + } + } + } Ok(()) } @@ -106,7 +124,7 @@ impl DlnaSession { } } -fn xml_escape(input: &str) -> String { +pub(crate) fn xml_escape(input: &str) -> String { let mut escaped = String::with_capacity(input.len()); for c in input.chars() { match c { @@ -120,3 +138,14 @@ fn xml_escape(input: &str) -> String { } escaped } + +#[cfg(test)] +mod tests { + use super::xml_escape; + + #[test] + fn xml_escape_covers_the_five_markup_chars() { + assert_eq!(xml_escape(r#"a&bd"e'f"#), "a&b<c>d"e'f"); + assert_eq!(xml_escape("http://10.0.0.1:1/t/playlist.m3u8"), "http://10.0.0.1:1/t/playlist.m3u8"); + } +} diff --git a/breadcast-core/src/http_server.rs b/breadcast-core/src/http_server.rs index 4565514..6907c07 100644 --- a/breadcast-core/src/http_server.rs +++ b/breadcast-core/src/http_server.rs @@ -16,8 +16,9 @@ const WORKER_THREADS: usize = 8; /// Serves `root` (the `hlssink3` output directory: `playlist.m3u8` + /// `segment*.ts`) over plain HTTP. Runs a small fixed pool of worker -/// threads; dropping the handle does not stop the server (there is no clean -/// shutdown yet — matches the smoke-testing scope of the rest of Phase 2). +/// threads. [`HttpServer::shutdown`] (also invoked from `Drop`) unblocks +/// those workers and closes the listener so a finished DLNA session does +/// not keep the last screen-recording segments reachable on the LAN. /// /// Every servable path is namespaced under a random token /// (`//playlist.m3u8`, etc. — see [`HttpServer::token`]) rather than @@ -37,6 +38,7 @@ const WORKER_THREADS: usize = 8; pub struct HttpServer { addr: SocketAddr, token: String, + server: Option>, } impl HttpServer { @@ -74,7 +76,19 @@ impl HttpServer { }); } - Ok(Self { addr, token }) + Ok(Self { addr, token, server: Some(server) }) + } + + /// Unblocks every worker and drops the listener. After this returns the + /// bind address is free and the path token no longer serves anything. + /// Safe to call more than once. + pub fn shutdown(&mut self) { + let Some(server) = self.server.take() else { return }; + // `unblock` wakes one `recv()` at a time; wake every worker so + // they all observe the error and drop their Arc. + for _ in 0..WORKER_THREADS { + server.unblock(); + } } /// The bound address, e.g. `0.0.0.0:41823`. Combine with this @@ -90,7 +104,17 @@ impl HttpServer { /// doc for why that can't just be `self.addr()`), including the /// unguessable path token every request must carry. pub fn url(&self, host: std::net::IpAddr, relative: &str) -> String { - format!("http://{host}:{}/{}/{relative}", self.addr.port(), self.token) + let port = self.addr.port(); + match host { + std::net::IpAddr::V6(v6) => format!("http://[{v6}]:{port}/{}/{relative}", self.token), + std::net::IpAddr::V4(v4) => format!("http://{v4}:{port}/{}/{relative}", self.token), + } + } +} + +impl Drop for HttpServer { + fn drop(&mut self) { + self.shutdown(); } } @@ -123,7 +147,7 @@ fn random_token() -> String { /// multi-range (multipart ranges aren't needed for HLS segment fetches, and /// falling back to a full 200 response for those is always a valid /// response under the HTTP spec). -fn parse_range(value: &str, len: usize) -> Option<(usize, usize)> { +pub(crate) fn parse_range(value: &str, len: usize) -> Option<(usize, usize)> { let spec = value.strip_prefix("bytes=")?; if spec.contains(',') || len == 0 { return None; @@ -233,7 +257,11 @@ fn handle_request(request: tiny_http::Request, root: &Path, token: &str) -> Resu ]; let (status, body) = match range { - Some((start, end)) if start <= end && end < data.len() => { + // RFC 7233: a last-byte-pos past the end is clamped, not 416. + // HLS clients often probe `bytes=0-1048575` against a ~400KB + // segment; answering 416 stalls playback with no encoder error. + Some((start, end)) if start < data.len() && start <= end => { + let end = end.min(data.len() - 1); headers.push(("Content-Range".to_string(), format!("bytes {start}-{end}/{}", data.len()))); (206u16, data[start..=end].to_vec()) } @@ -264,3 +292,25 @@ fn respond_status(request: tiny_http::Request, status: u16) -> Result<()> { .respond(tiny_http::Response::empty(status)) .context("failed to write HTTP error response") } + +#[cfg(test)] +mod tests { + use super::parse_range; + + #[test] + fn parse_range_accepts_the_usual_hls_shapes() { + assert_eq!(parse_range("bytes=0-99", 200), Some((0, 99))); + assert_eq!(parse_range("bytes=50-", 200), Some((50, 199))); + assert_eq!(parse_range("bytes=-20", 200), Some((180, 199))); + assert_eq!(parse_range("bytes=0-0", 200), Some((0, 0))); + } + + #[test] + fn parse_range_rejects_malformed_or_multipart() { + assert_eq!(parse_range("bytes=", 200), None); + assert_eq!(parse_range("bytes=-", 200), None); + assert_eq!(parse_range("bytes=0-10,20-30", 200), None); + assert_eq!(parse_range("items=0-10", 200), None); + assert_eq!(parse_range("bytes=0-10", 0), None); + } +} diff --git a/breadcast-core/src/net.rs b/breadcast-core/src/net.rs index fe20e75..4c74838 100644 --- a/breadcast-core/src/net.rs +++ b/breadcast-core/src/net.rs @@ -2,6 +2,8 @@ use std::net::{IpAddr, Ipv4Addr}; use anyhow::{Context, Result}; +const EXCLUDED_PREFIXES: &[&str] = &["tailscale", "wg", "docker", "veth", "br-", "virbr", "lo"]; + /// Finds this machine's LAN-reachable IPv4 address by enumerating network /// interfaces directly, rather than the more common "UDP-connect to a /// public address and read back the local endpoint" trick — that trick @@ -18,8 +20,42 @@ use anyhow::{Context, Result}; /// Shared by every casting protocol (Cast, DLNA, ...) — they all need to /// embed this machine's own address in a URL handed to a receiver device. pub fn local_lan_ip() -> Result { - const EXCLUDED_PREFIXES: &[&str] = &["tailscale", "wg", "docker", "veth", "br-", "virbr", "lo"]; + pick_lan_ip(None) +} +/// Like [`local_lan_ip`], but prefers the interface whose subnet contains +/// `peer`. Dual-homed machines (ethernet + wifi, two VLANs) otherwise +/// embed the wrong host in the HLS URL and the renderer cannot fetch it — +/// the same silent `LOAD FAILED` class as the Tailscale case above. +pub fn local_lan_ip_for(peer: IpAddr) -> Result { + pick_lan_ip(Some(peer)) +} + +fn pick_lan_ip(peer: Option) -> Result { + let mut candidates = lan_ipv4_candidates()?; + + if let Some(IpAddr::V4(peer_v4)) = peer { + if let Some((_, addr, _)) = candidates + .iter() + .filter(|(_, addr, prefix)| same_subnet(*addr, peer_v4, *prefix)) + .max_by_key(|(_, _, prefix)| *prefix) + { + return Ok(IpAddr::V4(*addr)); + } + } + + // Prefer a conventionally-named physical/Wi-Fi interface when there's a + // choice, but any private, non-excluded address is acceptable. + candidates.sort_by_key(|(iface, _, _)| !(iface.starts_with("wl") || iface.starts_with("en") || iface.starts_with("eth"))); + + candidates + .into_iter() + .map(|(_, addr, _)| IpAddr::V4(addr)) + .next() + .context("no LAN-reachable IPv4 address found (excluding loopback/VPN/virtual interfaces) — is this machine connected to a network?") +} + +fn lan_ipv4_candidates() -> Result> { let output = std::process::Command::new("ip") .args(["-4", "-o", "addr", "show", "scope", "global", "up"]) .output() @@ -29,7 +65,7 @@ pub fn local_lan_ip() -> Result { } let text = String::from_utf8_lossy(&output.stdout); - let mut candidates: Vec<(String, Ipv4Addr)> = Vec::new(); + let mut candidates: Vec<(String, Ipv4Addr, u8)> = Vec::new(); for line in text.lines() { // Format: "3: wlan0 inet 10.179.161.89/23 brd ... scope global dynamic wlan0" let mut fields = line.split_whitespace(); @@ -42,20 +78,41 @@ pub fn local_lan_ip() -> Result { continue; } let Some(cidr) = fields.next() else { continue }; - let Some(addr) = cidr.split('/').next().and_then(|a| a.parse::().ok()) else { continue }; + let mut parts = cidr.split('/'); + let Some(addr) = parts.next().and_then(|a| a.parse::().ok()) else { continue }; + let prefix: u8 = parts.next().and_then(|p| p.parse().ok()).unwrap_or(32); if !addr.is_private() { continue; } - candidates.push((iface.to_string(), addr)); + candidates.push((iface.to_string(), addr, prefix.min(32))); + } + Ok(candidates) +} + +fn same_subnet(a: Ipv4Addr, b: Ipv4Addr, prefix: u8) -> bool { + if prefix == 0 { + return true; + } + let mask = if prefix >= 32 { + u32::MAX + } else { + !((1u32 << (32 - prefix)) - 1) + }; + (u32::from(a) & mask) == (u32::from(b) & mask) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn same_subnet_respects_prefix_length() { + let a: Ipv4Addr = "10.179.161.89".parse().unwrap(); + let b: Ipv4Addr = "10.179.160.1".parse().unwrap(); + let other: Ipv4Addr = "10.0.0.1".parse().unwrap(); + assert!(same_subnet(a, b, 23)); + assert!(!same_subnet(a, other, 23)); + assert!(same_subnet(a, a, 32)); + assert!(!same_subnet(a, b, 32)); } - - // Prefer a conventionally-named physical/Wi-Fi interface when there's a - // choice, but any private, non-excluded address is acceptable. - candidates.sort_by_key(|(iface, _)| !(iface.starts_with("wl") || iface.starts_with("en") || iface.starts_with("eth"))); - - candidates - .into_iter() - .map(|(_, addr)| IpAddr::V4(addr)) - .next() - .context("no LAN-reachable IPv4 address found (excluding loopback/VPN/virtual interfaces) — is this machine connected to a network?") } diff --git a/breadcast-core/src/pipeline/mod.rs b/breadcast-core/src/pipeline/mod.rs index fe46558..ced7f4e 100644 --- a/breadcast-core/src/pipeline/mod.rs +++ b/breadcast-core/src/pipeline/mod.rs @@ -451,6 +451,7 @@ pub fn pull_encoded_frame(appsink: &gst_app::AppSink) -> Result, let poll = gst::ClockTime::from_mseconds(CAPTURE_STALL_POLL.as_millis() as u64); let mut stalled_for = Duration::ZERO; loop { + let pulled_at = std::time::Instant::now(); let Some(sample) = appsink.try_pull_sample(Some(poll)) else { // EOS is the ordinary end: the user hit "Stop sharing" in the // portal, or the source went away. @@ -458,16 +459,18 @@ pub fn pull_encoded_frame(appsink: &gst_app::AppSink) -> Result, return Ok(None); } // Teardown from another thread (`CastMirrorSession::stop` sets - // the pipeline to Null) makes the sink flush, and a flushing - // sink returns `None` *immediately* rather than after the - // timeout. Treat that as a clean end too -- otherwise this would - // busy-spin for the whole stall budget and then report a - // spurious "capture stalled" on every normal stop. - if !matches!(appsink.current_state(), gst::State::Playing | gst::State::Paused) { + // the pipeline to Null) flushes the sink *before* + // `current_state()` leaves Playing. A flushing sink returns + // `None` immediately; counting that as stall time used to + // raise a spurious "capture stalled" on every normal stop. + if !matches!(appsink.current_state(), gst::State::Playing | gst::State::Paused) + || matches!(appsink.pending_state(), gst::State::Null | gst::State::Ready) + || pulled_at.elapsed() < Duration::from_millis(20) + { return Ok(None); } - stalled_for += CAPTURE_STALL_POLL; + stalled_for += pulled_at.elapsed(); if stalled_for < CAPTURE_STALL_TIMEOUT { continue; } @@ -529,11 +532,56 @@ pub enum RunOutcome { Timeout, } +/// Blocks the calling thread until the pipeline reports an error or EOS. +/// Unlike [`run_until_error_or_timeout`] this does not give up after a +/// fixed duration -- a live mirror session can last hours, and a 1-hour +/// leftover from the smoke-test helper was leaving GStreamer errors +/// unobserved for the rest of the cast. Also returns [`RunOutcome::Eos`] +/// once the pipeline has been torn down from another thread (`Null`), so +/// a daemon watcher does not sit forever after `stop()`. +pub fn run_until_eos_or_error(pipeline: &gst::Pipeline) -> Result { + let bus = pipeline.bus().context("pipeline has no bus")?; + loop { + let Some(msg) = bus.timed_pop_filtered( + gst::ClockTime::from_mseconds(500), + &[gst::MessageType::Error, gst::MessageType::Eos, gst::MessageType::Warning], + ) else { + if !matches!( + pipeline.current_state(), + gst::State::Playing | gst::State::Paused | gst::State::Ready + ) { + return Ok(RunOutcome::Eos); + } + continue; + }; + + use gst::MessageView; + match msg.view() { + MessageView::Error(e) => { + bail!( + "GStreamer pipeline error from {:?}: {} ({:?})", + e.src().map(|s| s.path_string()), + e.error(), + e.debug() + ); + } + MessageView::Warning(w) => { + tracing::warn!( + src = ?w.src().map(|s| s.path_string()), + error = %w.error(), + "GStreamer pipeline warning" + ); + } + MessageView::Eos(_) => return Ok(RunOutcome::Eos), + _ => {} + } + } +} + /// Blocks the calling thread until the pipeline reports an error or EOS, or /// `timeout` elapses (whichever first). Returns which of those happened, or /// `Err` on a real pipeline error. Meant for smoke-testing from a -/// synchronous `main`/example; the real daemon will want an async/watch-based -/// version instead of blocking a thread. +/// synchronous `main`/example; the daemon uses [`run_until_eos_or_error`]. pub fn run_until_error_or_timeout(pipeline: &gst::Pipeline, timeout: gst::ClockTime) -> Result { let bus = pipeline.bus().context("pipeline has no bus")?; let deadline = std::time::Instant::now() + std::time::Duration::from(timeout); diff --git a/breadcast/src/css.rs b/breadcast/src/css.rs index 0f94064..73b7cb9 100644 --- a/breadcast/src/css.rs +++ b/breadcast/src/css.rs @@ -47,6 +47,11 @@ pub fn build_css(palette: &Palette) -> String { color: {overlay}; padding: 24px 8px; }} + .cast-error-label {{ + color: #e35b5b; + font-size: 0.85em; + padding: 4px 8px; + }} .cast-stop-button {{ background-color: alpha(#e35b5b, 0.15); color: #e35b5b; diff --git a/breadcast/src/ipc_client.rs b/breadcast/src/ipc_client.rs index 3b1ccf3..64250bb 100644 --- a/breadcast/src/ipc_client.rs +++ b/breadcast/src/ipc_client.rs @@ -29,6 +29,8 @@ impl IpcClient { let path = socket_path()?; let write_half = UnixStream::connect(&path) .with_context(|| format!("failed to connect to breadcastd at {} — is it running?", path.display()))?; + // A wedged daemon must not freeze the GTK main thread on write. + let _ = write_half.set_write_timeout(Some(std::time::Duration::from_secs(5))); let read_half = write_half.try_clone().context("failed to duplicate the IPC socket handle")?; let (tx, rx) = mpsc::channel(); diff --git a/breadcast/src/main.rs b/breadcast/src/main.rs index 85e9e02..56202d6 100644 --- a/breadcast/src/main.rs +++ b/breadcast/src/main.rs @@ -13,7 +13,9 @@ use std::rc::Rc; use bread_theme::load_palette; use breadcast_core::ipc::{DeviceInfo, Protocol, ServerMessage, StateInfo}; use gtk4::prelude::*; -use gtk4::{Align, Application, Box as GBox, Button, Label, ListBox, Orientation, SelectionMode, glib}; +use gtk4::{ + Align, Application, Box as GBox, Button, EventControllerKey, Label, ListBox, Orientation, SelectionMode, glib, +}; use ipc_client::IpcClient; const PANEL_WIDTH: i32 = 360; @@ -70,12 +72,50 @@ fn build_ui(app: &Application) { empty_label.add_css_class("cast-empty-label"); panel.append(&empty_label); + let error_label = Label::new(None); + error_label.add_css_class("cast-error-label"); + error_label.set_wrap(true); + error_label.set_visible(false); + panel.append(&error_label); + window.set_child(Some(&panel)); bread_utils::gtk_popup::close_on_outside_click(&window, &panel, { let window = window.clone(); move || window.close() }); + let key_ctrl = EventControllerKey::new(); + key_ctrl.set_propagation_phase(gtk4::PropagationPhase::Capture); + { + let window = window.clone(); + let list = list.clone(); + key_ctrl.connect_key_pressed(move |_, key, _, _| { + use gtk4::gdk::Key; + match key { + Key::Escape => { + window.close(); + glib::Propagation::Stop + } + Key::Return | Key::KP_Enter => { + if let Some(row) = list.selected_row() { + row.activate(); + } + glib::Propagation::Stop + } + Key::Down => { + bread_utils::gtk_popup::select_next_visible(&list); + glib::Propagation::Stop + } + Key::Up => { + bread_utils::gtk_popup::select_prev_visible(&list); + glib::Propagation::Stop + } + _ => glib::Propagation::Proceed, + } + }); + } + window.add_controller(key_ctrl); + match IpcClient::connect() { Ok((client, rx)) => { let client = Rc::new(RefCell::new(client)); @@ -84,10 +124,16 @@ fn build_ui(app: &Application) { list.connect_row_activated({ let client = client.clone(); + let window = window.clone(); move |_, row| { let Some(device_id) = (unsafe { row.data::("device_id") }) else { return }; let device_id = unsafe { device_id.as_ref() }.clone(); let _ = client.borrow_mut().send("start_cast", serde_json::json!({ "device_id": device_id })); + // Close before the portal picker appears -- this overlay + // is KeyboardMode::Exclusive and otherwise steals + // Enter/Escape from the picker (and swallows clicks on + // the dimmed background). + window.close(); } }); @@ -100,11 +146,26 @@ fn build_ui(app: &Application) { let list = list.clone(); let empty_label = empty_label.clone(); + let error_label = error_label.clone(); let status_pill = status_pill.clone(); let stop_button = stop_button.clone(); glib::timeout_add_local(std::time::Duration::from_millis(100), move || { - while let Ok(message) = rx.try_recv() { - handle_server_message(message, &list, &empty_label, &status_pill, &stop_button); + loop { + match rx.try_recv() { + Ok(message) => { + handle_server_message(message, &list, &empty_label, &error_label, &status_pill, &stop_button); + } + Err(std::sync::mpsc::TryRecvError::Empty) => break, + Err(std::sync::mpsc::TryRecvError::Disconnected) => { + empty_label.set_label("breadcastd isn't running"); + empty_label.set_visible(true); + list.set_visible(false); + stop_button.set_visible(false); + status_pill.set_label("Offline"); + status_pill.remove_css_class("casting"); + return glib::ControlFlow::Break; + } + } } glib::ControlFlow::Continue }); @@ -125,6 +186,7 @@ fn handle_server_message( message: ServerMessage, list: &ListBox, empty_label: &Label, + error_label: &Label, status_pill: &Label, stop_button: &Button, ) { @@ -138,6 +200,8 @@ fn handle_server_message( } ServerMessage::Response { error: Some(error), .. } => { eprintln!("breadcast: request failed: {error}"); + error_label.set_label(&error); + error_label.set_visible(true); None } _ => None, @@ -145,7 +209,10 @@ fn handle_server_message( let Some((event, data)) = payload else { return }; match event.as_str() { "device_list_changed" => update_device_list(list, empty_label, data), - "state_changed" => update_state(status_pill, stop_button, data), + "state_changed" => { + error_label.set_visible(false); + update_state(status_pill, stop_button, data); + } _ => {} } } diff --git a/breadcastd/src/cast_mirror.rs b/breadcastd/src/cast_mirror.rs index 2990077..b8444ad 100644 --- a/breadcastd/src/cast_mirror.rs +++ b/breadcastd/src/cast_mirror.rs @@ -83,7 +83,11 @@ impl CastMirrorSession { /// receiver dropping the connection) back to the daemon actor, so it /// can transition back to `Idle` and notify GUI clients even if nobody /// called `stop()`. - pub async fn start(device: CastDevice, daemon_tx: tokio::sync::mpsc::Sender) -> Result { + pub async fn start( + device: CastDevice, + daemon_tx: tokio::sync::mpsc::Sender, + generation: u64, + ) -> Result { let capture = CaptureSession::start().await.context("failed to start portal screen capture")?; let video_node_id = capture.video_node_id(); @@ -94,13 +98,18 @@ impl CastMirrorSession { // OFFER below rather than re-derived there, so the advertised stream // and the encoded stream cannot drift apart. let (pipeline, appsink, encoder, video_params) = - build_video_pipeline_for_streaming(video_node_id).context("failed to build the encode pipeline")?; + match build_video_pipeline_for_streaming(video_node_id) { + Ok(built) => built, + Err(e) => { + close_capture_bounded(capture).await; + return Err(e).context("failed to build the encode pipeline"); + } + }; { let pipeline_watch = pipeline.clone(); std::thread::spawn(move || { - match breadcast_core::pipeline::run_until_error_or_timeout(&pipeline_watch, gst::ClockTime::from_seconds(3600)) - { + match breadcast_core::pipeline::run_until_eos_or_error(&pipeline_watch) { Ok(outcome) => tracing::debug!(?outcome, "encode pipeline bus watcher ended"), Err(e) => tracing::error!(error = ?e, "encode pipeline error"), } @@ -111,19 +120,37 @@ impl CastMirrorSession { // (milliseconds on a LAN) but still blocking I/O -- run it off the // async worker thread pool rather than stalling it, even briefly. let device_for_connect = device.clone(); - let (session, _media_events, raw_messages) = tokio::task::spawn_blocking(move || { + let connect = tokio::task::spawn_blocking(move || { CastSession::connect_app( &device_for_connect, CastDeviceApp::Custom(breadcast_core::caststream::MIRRORING_APP_ID.to_string()), ) }) - .await - .context("connect_app task panicked")? - .context("failed to connect and launch the Mirroring receiver")?; + .await; + let (session, _media_events, raw_messages) = match connect { + Ok(Ok(connected)) => connected, + Ok(Err(e)) => { + let _ = pipeline.set_state(gst::State::Null); + close_capture_bounded(capture).await; + return Err(e).context("failed to connect and launch the Mirroring receiver"); + } + Err(e) => { + let _ = pipeline.set_state(gst::State::Null); + close_capture_bounded(capture).await; + return Err(e).context("connect_app task panicked"); + } + }; let (sender, stream_events) = - CastStreamSender::start(&device.host, "sender-0", session.transport_id(), video_params) - .context("failed to start the Cast Streaming session")?; + match CastStreamSender::start(&device.host, "sender-0", session.transport_id(), video_params) { + Ok(started) => started, + Err(e) => { + stop_session_bounded(&session).await; + let _ = pipeline.set_state(gst::State::Null); + close_capture_bounded(capture).await; + return Err(e).context("failed to start the Cast Streaming session"); + } + }; let sender = Arc::new(sender); let message_pump = { @@ -138,9 +165,12 @@ impl CastMirrorSession { }; let negotiated = Arc::new(AtomicBool::new(false)); + let failed = Arc::new(AtomicBool::new(false)); let event_pump = { let session = session.clone(); let negotiated = negotiated.clone(); + let failed = failed.clone(); + let daemon_tx = daemon_tx.clone(); std::thread::spawn(move || { while let Ok(event) = stream_events.recv() { match event { @@ -150,37 +180,68 @@ impl CastMirrorSession { } } CastStreamEvent::Negotiated => negotiated.store(true, Ordering::Release), - CastStreamEvent::Error(message) => tracing::warn!(%message, "Cast Streaming error"), + CastStreamEvent::Error(message) => { + tracing::warn!(%message, "Cast Streaming error"); + failed.store(true, Ordering::Release); + // After negotiation the frame pump is running + // and this is an unprompted death; before + // negotiation, start() itself observes `failed` + // and tears down. + if negotiated.load(Ordering::Acquire) { + let _ = daemon_tx.blocking_send(DaemonCommand::SessionEnded { generation }); + } + } CastStreamEvent::PictureLost => tracing::debug!("receiver reported picture loss"), } } }) }; + // From here every error path must use the same join/drop order as + // `stop()`. Building the session now and calling `stop()` on it is + // what keeps a leaked `CastStreamSender` from leaving + // PlatformClientPosix alive -- the next start would then hit + // OSP_CHECK(!instance_) and abort the daemon. + let mut started = Self { + pipeline, + session, + capture: Some(capture), + sender: Some(sender), + message_pump: Some(message_pump), + event_pump: Some(event_pump), + frame_pump: None, + }; + tracing::info!(device = %device.name, "sending Cast Streaming OFFER"); - sender.negotiate(); + if let Some(sender) = started.sender.as_ref() { + sender.negotiate(); + } let deadline = tokio::time::Instant::now() + tokio::time::Duration::from_secs(10); - while !negotiated.load(Ordering::Acquire) && tokio::time::Instant::now() < deadline { + while !negotiated.load(Ordering::Acquire) + && !failed.load(Ordering::Acquire) + && tokio::time::Instant::now() < deadline + { tokio::time::sleep(tokio::time::Duration::from_millis(50)).await; } + if failed.load(Ordering::Acquire) { + started.stop().await; + anyhow::bail!("Cast Streaming session error from {} during negotiation", device.name); + } if !negotiated.load(Ordering::Acquire) { - // Bounded for the same reason `Self::stop`'s calls are -- an - // unresponsive receiver (which is exactly what "negotiation - // timed out" implies) can wedge either of these forever - // otherwise, taking the whole single-threaded daemon actor - // down with it before this even gets to return an error. - stop_session_bounded(&session).await; - close_capture_bounded(capture).await; + started.stop().await; anyhow::bail!("never received an ANSWER from {} (negotiation timed out)", device.name); } - pipeline.set_state(gst::State::Playing).context("failed to start the encode pipeline")?; + if let Err(e) = started.pipeline.set_state(gst::State::Playing) { + started.stop().await; + return Err(e).context("failed to start the encode pipeline"); + } tracing::info!(device = %device.name, "mirroring started"); let frame_pump = { let device_name = device.name.clone(); - let sender = sender.clone(); + let sender = started.sender.as_ref().expect("sender installed above").clone(); std::thread::spawn(move || { let result = frame_pump_loop(&appsink, &encoder, &sender); if let Err(e) = result { @@ -190,19 +251,12 @@ impl CastMirrorSession { // was, very recently) still alive. If the channel is full or // closed, there's nothing more useful to do from this // thread than drop the notification. - let _ = daemon_tx.blocking_send(DaemonCommand::SessionEnded); + let _ = daemon_tx.blocking_send(DaemonCommand::SessionEnded { generation }); }) }; + started.frame_pump = Some(frame_pump); - Ok(Self { - pipeline, - session, - capture: Some(capture), - sender: Some(sender), - message_pump: Some(message_pump), - event_pump: Some(event_pump), - frame_pump: Some(frame_pump), - }) + Ok(started) } /// Tears down the session. Order matters and is not interchangeable: diff --git a/breadcastd/src/daemon.rs b/breadcastd/src/daemon.rs index 95a0d30..1a164bb 100644 --- a/breadcastd/src/daemon.rs +++ b/breadcastd/src/daemon.rs @@ -33,7 +33,29 @@ pub enum DaemonCommand { /// dropping the connection or stopping playback) -- as opposed to /// `StopCast` being called. Either way the daemon needs to forget the /// (now-dead) session and go back to `Idle`. - SessionEnded, + SessionEnded { generation: u64 }, + /// Result of a `StartCast` that ran off this actor (the portal picker + /// and OFFER/ANSWER wait must not stall `list_devices` / `stop_cast`). + /// Ignored when `generation` no longer matches -- that means `StopCast` + /// cancelled the in-flight start. + StartFinished { + generation: u64, + outcome: Result, + reply: oneshot::Sender>, + }, +} + +/// A session that finished starting, ready to be installed as `active_session`. +pub(crate) struct StartedSession { + pub session: ActiveSession, + pub device_id: String, + pub device_name: String, + pub protocol: Protocol, +} + +pub(crate) struct StartFailed { + pub device_id: String, + pub error: String, } pub fn spawn(events_tx: broadcast::Sender, bread_client: BreadClient) -> mpsc::Sender { @@ -45,6 +67,8 @@ pub fn spawn(events_tx: broadcast::Sender, bread_client: BreadCli dlna_devices: HashMap::new(), state: StateInfo::Idle, active_session: None, + starting: false, + start_generation: 0, events_tx, bread_client, self_tx, @@ -59,7 +83,7 @@ pub fn spawn(events_tx: broadcast::Sender, bread_client: BreadCli /// The currently active mirroring session, if any -- exactly one of the two /// protocol-specific session types, chosen by which device map `StartCast` /// found the requested device id in. -enum ActiveSession { +pub(crate) enum ActiveSession { Cast(CastMirrorSession), Dlna(Box), } @@ -86,6 +110,13 @@ struct Daemon { dlna_devices: HashMap, state: StateInfo, active_session: Option, + /// True while a `StartCast` is running off this actor (portal picker / + /// negotiation). Distinct from `active_session` so a second start is + /// rejected before the first one has a session to install. + starting: bool, + /// Bumped by `StopCast` so a `StartFinished` from a cancelled start + /// tears its session down instead of installing it. + start_generation: u64, events_tx: broadcast::Sender, /// Used to publish `bread.cast.mirroring_started`/`.stopped`/`.failed` /// on state transitions -- see `bread_events.rs`. A no-op if breadd @@ -111,6 +142,12 @@ impl Daemon { self.start_cast(device_id, reply).await; } DaemonCommand::StopCast { reply } => { + if self.starting { + // Invalidate the in-flight start so its StartFinished + // tears the session down instead of installing it. + self.starting = false; + self.start_generation = self.start_generation.wrapping_add(1); + } if let Some(session) = self.active_session.take() { session.stop().await; bread_events::emit_mirroring_stopped(&self.bread_client); @@ -119,20 +156,32 @@ impl Daemon { self.broadcast_state(); let _ = reply.send(Ok(())); } - DaemonCommand::SessionEnded => { - // The session already tore itself down (that's what - // triggered this) -- just drop our handle to it and update - // state. Ignored if this arrives after an explicit - // `StopCast` already cleared `active_session` (the - // background pump/poll it came from may briefly outlive - // that call). - if self.active_session.take().is_some() { + DaemonCommand::SessionEnded { generation } => { + // Ignore a late notification from a session that already + // stopped (or whose start was cancelled). An abandoned + // pump after a 5s join timeout used to take() the *next* + // session and flip the UI to Idle while it was still live. + if generation != self.start_generation { + return; + } + // The session's pump/poll noticed death, but the session + // handle itself has *not* been torn down -- there is no + // Drop impl. Dropping it here used to skip CastSession::stop + // (TV left on the last frame), skip portal close (PipeWire + // leak), and destroy the FFI sender while pump threads were + // still calling into it. Always run the ordered stop. + if let Some(session) = self.active_session.take() { tracing::info!("mirror session ended on its own, returning to idle"); + session.stop().await; + self.start_generation = self.start_generation.wrapping_add(1); self.state = StateInfo::Idle; self.broadcast_state(); bread_events::emit_mirroring_stopped(&self.bread_client); } } + DaemonCommand::StartFinished { generation, outcome, reply } => { + self.on_start_finished(generation, outcome, reply).await; + } DaemonCommand::CastDeviceFound(device) => { // mDNS resolves one physical device on every local address // it has -- typically a private IPv4 and a link-local IPv6 @@ -147,8 +196,8 @@ impl Daemon { ); if should_replace { self.cast_devices.insert(device.id.clone(), device); + self.broadcast_devices(); } - self.broadcast_devices(); } DaemonCommand::CastDeviceLost(id) => { self.cast_devices.remove(&id); @@ -166,54 +215,100 @@ impl Daemon { } async fn start_cast(&mut self, device_id: String, reply: oneshot::Sender>) { - if self.active_session.is_some() { + if self.active_session.is_some() || self.starting { let _ = reply.send(Err("already casting -- stop the current session first".to_string())); return; } if let Some(device) = self.cast_devices.get(&device_id).cloned() { - match CastMirrorSession::start(device.clone(), self.self_tx.clone()).await { - Ok(session) => { - self.active_session = Some(ActiveSession::Cast(session)); - self.state = - StateInfo::Casting { device_id: device.id.clone(), device_name: device.name.clone(), protocol: Protocol::Cast }; - self.broadcast_state(); - bread_events::emit_mirroring_started(&self.bread_client, &device.id, &device.name, "cast"); - let _ = reply.send(Ok(())); - } - Err(e) => { + self.starting = true; + let generation = self.start_generation; + let daemon_tx = self.self_tx.clone(); + tokio::spawn(async move { + let outcome = match CastMirrorSession::start(device.clone(), daemon_tx.clone(), generation).await { + Ok(session) => Ok(StartedSession { + session: ActiveSession::Cast(session), + device_id: device.id, + device_name: device.name, + protocol: Protocol::Cast, + }), // `{e:#}` (not `{e}`/`to_string()`) so the full anyhow // context chain reaches the caller/GUI instead of just // the outermost ".context()" message. - bread_events::emit_mirroring_failed(&self.bread_client, &device.id, &format!("{e:#}")); - let _ = reply.send(Err(format!("{e:#}"))); - } - } + Err(e) => Err(StartFailed { device_id: device.id, error: format!("{e:#}") }), + }; + let _ = daemon_tx.send(DaemonCommand::StartFinished { generation, outcome, reply }).await; + }); return; } if let Some(device) = self.dlna_devices.get(&device_id).cloned() { - match DlnaMirrorSession::start(device.clone(), self.self_tx.clone()).await { - Ok(session) => { - self.active_session = Some(ActiveSession::Dlna(Box::new(session))); - self.state = StateInfo::Casting { - device_id: device.url.clone(), - device_name: device.friendly_name.clone(), + self.starting = true; + let generation = self.start_generation; + let daemon_tx = self.self_tx.clone(); + tokio::spawn(async move { + let outcome = match DlnaMirrorSession::start(device.clone(), daemon_tx.clone(), generation).await { + Ok(session) => Ok(StartedSession { + session: ActiveSession::Dlna(Box::new(session)), + device_id: device.url, + device_name: device.friendly_name, protocol: Protocol::Dlna, - }; - self.broadcast_state(); - bread_events::emit_mirroring_started(&self.bread_client, &device.url, &device.friendly_name, "dlna"); - let _ = reply.send(Ok(())); - } - Err(e) => { - bread_events::emit_mirroring_failed(&self.bread_client, &device.url, &format!("{e:#}")); - let _ = reply.send(Err(format!("{e:#}"))); - } - } + }), + Err(e) => Err(StartFailed { device_id: device.url, error: format!("{e:#}") }), + }; + let _ = daemon_tx.send(DaemonCommand::StartFinished { generation, outcome, reply }).await; + }); return; } - let _ = reply.send(Err(format!("unknown device id \"{device_id}\""))); + let error = format!("unknown device id \"{device_id}\""); + bread_events::emit_mirroring_failed(&self.bread_client, &device_id, &error); + let _ = reply.send(Err(error)); + } + + async fn on_start_finished( + &mut self, + generation: u64, + outcome: Result, + reply: oneshot::Sender>, + ) { + if generation != self.start_generation || !self.starting { + // StopCast cancelled this start while the portal/negotiation + // was still running. Tear the session down if it succeeded + // anyway, so we don't leak a sender or leave the TV casting. + if let Ok(started) = outcome { + started.session.stop().await; + } + let _ = reply.send(Err("start cancelled".to_string())); + return; + } + self.starting = false; + match outcome { + Ok(started) => { + self.active_session = Some(started.session); + self.state = StateInfo::Casting { + device_id: started.device_id.clone(), + device_name: started.device_name.clone(), + protocol: started.protocol, + }; + self.broadcast_state(); + let protocol = match started.protocol { + Protocol::Cast => "cast", + Protocol::Dlna => "dlna", + }; + bread_events::emit_mirroring_started( + &self.bread_client, + &started.device_id, + &started.device_name, + protocol, + ); + let _ = reply.send(Ok(())); + } + Err(failed) => { + bread_events::emit_mirroring_failed(&self.bread_client, &failed.device_id, &failed.error); + let _ = reply.send(Err(failed.error)); + } + } } fn device_list(&self) -> Vec { diff --git a/breadcastd/src/dlna_mirror.rs b/breadcastd/src/dlna_mirror.rs index c3c68e2..510cc00 100644 --- a/breadcastd/src/dlna_mirror.rs +++ b/breadcastd/src/dlna_mirror.rs @@ -13,8 +13,11 @@ use std::time::Duration; use anyhow::{Context, Result}; +use breadcast_core::http_server::HttpServer; use breadcast_core::net::local_lan_ip; -use breadcast_core::pipeline::{build_video_pipeline, hls_output_dir, wait_for_playlist_segments}; +use breadcast_core::pipeline::{ + build_video_pipeline, hls_output_dir, run_until_eos_or_error, wait_for_playlist_segments, +}; use breadcast_core::{CaptureSession, DlnaDevice, DlnaSession}; use gstreamer as gst; use gstreamer::prelude::*; @@ -27,11 +30,20 @@ use crate::daemon::DaemonCommand; /// this is a poll, not a push. const POLL_INTERVAL: Duration = Duration::from_secs(3); +/// How long the playlist may sit unchanged before we treat the encode +/// side as dead. The Cast path has `pull_encoded_frame`'s stall watchdog; +/// this is the HLS equivalent -- portal/encoder stalls produce no bus +/// error, just a playlist that stops growing. +const PLAYLIST_STALL_TIMEOUT: Duration = Duration::from_secs(15); + pub struct DlnaMirrorSession { pipeline: gst::Pipeline, session: DlnaSession, capture: Option, + http: HttpServer, + output_dir: std::path::PathBuf, poll_task: tokio::task::JoinHandle<()>, + stall_task: tokio::task::JoinHandle<()>, } impl DlnaMirrorSession { @@ -43,42 +55,57 @@ impl DlnaMirrorSession { /// renderer stopping playback on its own, a GStreamer error, or the /// renderer becoming unreachable) back to the daemon actor — mirrors /// `CastMirrorSession::start`'s same use of it. - pub async fn start(device: DlnaDevice, daemon_tx: tokio::sync::mpsc::Sender) -> Result { + pub async fn start( + device: DlnaDevice, + daemon_tx: tokio::sync::mpsc::Sender, + generation: u64, + ) -> Result { let capture = CaptureSession::start().await.context("failed to start portal screen capture")?; let video_node_id = capture.video_node_id(); - let output_dir = hls_output_dir("dlna-mirror")?; - let pipeline = - build_video_pipeline(video_node_id, &output_dir).context("failed to build the encode pipeline")?; + let output_dir = match hls_output_dir("dlna-mirror") { + Ok(dir) => dir, + Err(e) => { + let _ = capture.close().await; + return Err(e); + } + }; + let pipeline = match build_video_pipeline(video_node_id, &output_dir) { + Ok(p) => p, + Err(e) => { + let _ = capture.close().await; + let _ = std::fs::remove_dir_all(&output_dir); + return Err(e).context("failed to build the encode pipeline"); + } + }; - // Fire-and-forget, same as `CastMirrorSession::start`'s identical - // block: nothing joins this thread, it just self-terminates on - // pipeline error, EOS, or its own 1-hour timeout, whichever is - // first — see that function's doc comment for why that's fine. - { - let pipeline_watch = pipeline.clone(); - std::thread::spawn(move || { - match breadcast_core::pipeline::run_until_error_or_timeout(&pipeline_watch, gst::ClockTime::from_seconds(3600)) - { - Ok(outcome) => tracing::debug!(?outcome, "DLNA encode pipeline bus watcher ended"), - Err(e) => tracing::error!(error = ?e, "DLNA encode pipeline error"), - } - }); + if let Err(e) = pipeline.set_state(gst::State::Playing) { + let _ = capture.close().await; + let _ = std::fs::remove_dir_all(&output_dir); + return Err(e).context("failed to start the encode pipeline"); } - pipeline.set_state(gst::State::Playing).context("failed to start the encode pipeline")?; - - let lan_ip = local_lan_ip().context("failed to determine this machine's LAN-reachable IP")?; + let peer = host_ip_from_url(&device.url); + let lan_ip = match peer.map(breadcast_core::net::local_lan_ip_for).unwrap_or_else(local_lan_ip) { + Ok(ip) => ip, + Err(e) => { + abort_partial(&pipeline, capture, None, &output_dir).await; + return Err(e).context("failed to determine this machine's LAN-reachable IP"); + } + }; // Bind an ephemeral port (`:0`) rather than a fixed one like the // `dlna_mirror_test` example uses -- the daemon may need to run // alongside that example, or a future concurrent-session mode, - // without a bind conflict. `HttpServer::start`'s worker threads - // outlive this session once it stops (documented pre-existing - // limitation, see `http_server.rs` -- not something introduced - // here); one leaked idle listener per DLNA cast is an accepted - // cost until that gets a real shutdown path. - let http = breadcast_core::http_server::HttpServer::start("0.0.0.0:0", output_dir.clone()) - .context("failed to start the HLS HTTP server")?; + // without a bind conflict. The server is held on the session and + // shut down in `stop()` so the last screen-recording segments are + // not left reachable on the LAN. + let http = match HttpServer::start("0.0.0.0:0", output_dir.clone()) { + Ok(http) => http, + Err(e) => { + abort_partial(&pipeline, capture, None, &output_dir).await; + return Err(e).context("failed to start the HLS HTTP server"); + } + }; let stream_url = http.url(lan_ip, "playlist.m3u8"); // Two segments, not three: this is a "don't hand the renderer a 404 @@ -87,30 +114,42 @@ impl DlnaMirrorSession { // `build_video_pipeline`'s note on HLS latency). Two is the minimum // that still proves the encoder is genuinely producing output rather // than having emitted one segment and stalled. - wait_for_playlist_segments(&output_dir.join("playlist.m3u8"), 2, Duration::from_secs(20)) - .await - .context("encode pipeline never produced playable HLS segments")?; + if let Err(e) = wait_for_playlist_segments(&output_dir.join("playlist.m3u8"), 2, Duration::from_secs(20)).await + { + abort_partial(&pipeline, capture, Some(http), &output_dir).await; + return Err(e).context("encode pipeline never produced playable HLS segments"); + } - let session = DlnaSession::connect(&device).await.context("failed to connect to the DLNA renderer")?; - session.load(&stream_url).await.context("renderer rejected the stream load")?; + let session = match DlnaSession::connect(&device).await { + Ok(session) => session, + Err(e) => { + abort_partial(&pipeline, capture, Some(http), &output_dir).await; + return Err(e).context("failed to connect to the DLNA renderer"); + } + }; + if let Err(e) = session.load(&stream_url).await { + abort_partial(&pipeline, capture, Some(http), &output_dir).await; + return Err(e).context("renderer rejected the stream load"); + } tracing::info!(device = %device.friendly_name, %stream_url, "DLNA mirroring started"); let poll_task = { let session = session.clone(); let device_name = device.friendly_name.clone(); + let daemon_tx = daemon_tx.clone(); tokio::spawn(async move { loop { tokio::time::sleep(POLL_INTERVAL).await; match session.transport_state().await { Ok(state) if state == "STOPPED" || state == "NO_MEDIA_PRESENT" => { tracing::info!(device = %device_name, %state, "DLNA renderer ended playback on its own"); - let _ = daemon_tx.send(DaemonCommand::SessionEnded).await; + let _ = daemon_tx.send(DaemonCommand::SessionEnded { generation }).await; return; } Ok(_) => {} Err(e) => { tracing::warn!(device = %device_name, error = ?e, "DLNA transport state poll failed, treating renderer as gone"); - let _ = daemon_tx.send(DaemonCommand::SessionEnded).await; + let _ = daemon_tx.send(DaemonCommand::SessionEnded { generation }).await; return; } } @@ -118,17 +157,69 @@ impl DlnaMirrorSession { }) }; - Ok(Self { pipeline, session, capture: Some(capture), poll_task }) + let stall_task = { + let playlist = output_dir.join("playlist.m3u8"); + let daemon_tx = daemon_tx.clone(); + tokio::spawn(async move { + let mut last_mtime = None; + let mut stalled_for = Duration::ZERO; + loop { + tokio::time::sleep(Duration::from_secs(1)).await; + let mtime = std::fs::metadata(&playlist).and_then(|m| m.modified()).ok(); + if mtime != last_mtime { + last_mtime = mtime; + stalled_for = Duration::ZERO; + continue; + } + stalled_for += Duration::from_secs(1); + if stalled_for >= PLAYLIST_STALL_TIMEOUT { + tracing::warn!( + "DLNA encode stalled: playlist unchanged for {}s", + PLAYLIST_STALL_TIMEOUT.as_secs() + ); + let _ = daemon_tx.send(DaemonCommand::SessionEnded { generation }).await; + return; + } + } + }) + }; + + // Portal "stop sharing" / a real GStreamer error used to only log + // -- the bus watcher never told the daemon, so the UI stayed on + // Casting and the renderer kept looping stale segments. Notify. + { + let pipeline_watch = pipeline.clone(); + let daemon_tx = daemon_tx.clone(); + std::thread::spawn(move || { + match run_until_eos_or_error(&pipeline_watch) { + Ok(outcome) => tracing::debug!(?outcome, "DLNA encode pipeline bus watcher ended"), + Err(e) => tracing::error!(error = ?e, "DLNA encode pipeline error"), + } + let _ = daemon_tx.blocking_send(DaemonCommand::SessionEnded { generation }); + }); + } + + Ok(Self { + pipeline, + session, + capture: Some(capture), + http, + output_dir, + poll_task, + stall_task, + }) } /// Tears down the session: stops polling, tells the renderer to stop, - /// stops the encode pipeline, and closes the portal capture session. + /// stops the encode pipeline, shuts the HLS server, closes the portal + /// capture session, and deletes the recording directory. pub async fn stop(mut self) { // A request, not a wait -- if the poll task is mid-poll and sends // one more `SessionEnded` right as this races it, that's harmless: // `daemon.rs`'s handler already no-ops when `active_session` was // already cleared by this explicit stop. self.poll_task.abort(); + self.stall_task.abort(); if let Err(e) = self.session.stop().await { tracing::warn!(error = ?e, "failed to cleanly stop the DLNA session"); @@ -136,10 +227,45 @@ impl DlnaMirrorSession { if let Err(e) = self.pipeline.set_state(gst::State::Null) { tracing::warn!(error = ?e, "failed to stop the encode pipeline cleanly"); } + self.http.shutdown(); if let Some(capture) = self.capture.take() { - if let Err(e) = capture.close().await { - tracing::warn!(error = ?e, "failed to cleanly close the portal capture session"); + match tokio::time::timeout(Duration::from_secs(5), capture.close()).await { + Ok(Err(e)) => tracing::warn!(error = ?e, "failed to cleanly close the portal capture session"), + Err(_) => tracing::warn!("portal capture session did not close within 5s -- abandoning it"), + Ok(Ok(())) => {} } } + if let Err(e) = std::fs::remove_dir_all(&self.output_dir) { + tracing::debug!(error = %e, dir = %self.output_dir.display(), "failed to remove HLS output dir"); + } } } + +async fn abort_partial( + pipeline: &gst::Pipeline, + capture: CaptureSession, + http: Option, + output_dir: &std::path::Path, +) { + let _ = pipeline.set_state(gst::State::Null); + if let Some(mut http) = http { + http.shutdown(); + } + match tokio::time::timeout(Duration::from_secs(5), capture.close()).await { + Ok(Err(e)) => tracing::warn!(error = ?e, "failed to close portal capture after a failed DLNA start"), + Err(_) => tracing::warn!("portal capture did not close within 5s after a failed DLNA start"), + Ok(Ok(())) => {} + } + let _ = std::fs::remove_dir_all(output_dir); +} + +fn host_ip_from_url(url: &str) -> Option { + let rest = url.split("://").nth(1)?; + let hostport = rest.split('/').next()?; + let host = if let Some(inside) = hostport.strip_prefix('[') { + inside.split(']').next()? + } else { + hostport.rsplit_once(':').map(|(h, _)| h).unwrap_or(hostport) + }; + host.parse().ok() +} diff --git a/breadcastd/src/ipc.rs b/breadcastd/src/ipc.rs index d8cbd85..703d50c 100644 --- a/breadcastd/src/ipc.rs +++ b/breadcastd/src/ipc.rs @@ -13,6 +13,8 @@ //! a time, but nothing here assumes that) can connect concurrently; each //! gets its own copy of every broadcast event. +use std::os::unix::fs::DirBuilderExt; + use anyhow::{Context, Result}; use breadcast_core::ipc::{ClientRequest, ServerMessage, socket_path}; use serde_json::Value; @@ -31,7 +33,13 @@ use crate::daemon::DaemonCommand; pub async fn serve(daemon_tx: mpsc::Sender, events_tx: broadcast::Sender) -> Result<()> { let socket_path = socket_path()?; if let Some(parent) = socket_path.parent() { - std::fs::create_dir_all(parent) + // 0700 even if umask is loose -- this directory holds the control + // socket, and XDG_RUNTIME_DIR itself is 0700 but a recreate after + // a wiped runtime dir should not inherit a world-readable mode. + std::fs::DirBuilder::new() + .recursive(true) + .mode(0o700) + .create(parent) .with_context(|| format!("failed to create socket dir {}", parent.display()))?; } // A stale socket file from an unclean previous exit makes bind() fail @@ -45,7 +53,18 @@ pub async fn serve(daemon_tx: mpsc::Sender, events_tx: broadcast: tracing::info!(path = %socket_path.display(), "IPC socket listening"); loop { - let (stream, _addr) = listener.accept().await.context("failed to accept IPC connection")?; + let (stream, _addr) = match listener.accept().await { + Ok(accepted) => accepted, + Err(e) => { + // EMFILE / a single bad accept must not take the control + // socket down for the rest of the daemon's life -- discovery + // would keep running while every `breadcast` launch reports + // "isn't running". + tracing::warn!(error = %e, "IPC accept failed, retrying"); + tokio::time::sleep(std::time::Duration::from_millis(100)).await; + continue; + } + }; let daemon_tx = daemon_tx.clone(); let events_rx = events_tx.subscribe(); tokio::spawn(async move { @@ -106,10 +125,14 @@ async fn handle_connection( continue; } }; - let response = handle_request(request, &daemon_tx).await; - if writer_tx.send(response).is_err() { - break; - } + // Handle off the read loop so a long `start_cast` (portal picker) + // does not block `stop_cast` sitting in the same socket buffer. + let daemon_tx = daemon_tx.clone(); + let writer_tx = writer_tx.clone(); + tokio::spawn(async move { + let response = handle_request(request, &daemon_tx).await; + let _ = writer_tx.send(response); + }); } forward_task.abort(); diff --git a/contrib/breadcastd.service b/contrib/breadcastd.service index c7778b5..69571af 100644 --- a/contrib/breadcastd.service +++ b/contrib/breadcastd.service @@ -7,7 +7,7 @@ PartOf=graphical-session.target [Service] Type=simple -ExecStart=%h/.cargo/bin/breadcastd +ExecStart=%h/.local/bin/breadcastd Restart=on-failure RestartSec=2 diff --git a/vendor/rust_cast-0.21.0/PATCHES.md b/vendor/rust_cast-0.21.0/PATCHES.md index 50f4b8f..94308d9 100644 --- a/vendor/rust_cast-0.21.0/PATCHES.md +++ b/vendor/rust_cast-0.21.0/PATCHES.md @@ -22,11 +22,14 @@ channel claims the namespace. ## The patch -`src/lib.rs`: added `CastDevice::send_message(&self, namespace, -destination, message)`, built the same way `ReceiverChannel::broadcast_message()` -is internally, but with a caller-supplied `destination` instead of a -hardcoded `"*"`. See the doc comment on that method for the exact rationale -(marked "LOCAL PATCH (breadcast, not upstream)"). +`src/lib.rs`: added `CastDevice::send_message(&self, namespace: &str, +destination: &str, message: &str)`, built the same way +`ReceiverChannel::broadcast_message()` is internally, but with a +caller-supplied `destination` instead of a hardcoded `"*"`. The payload is +a raw string (already-serialized JSON from the vendored openscreen C++), +not a `Serialize` value -- a generic `M: Serialize` would double-encode +the OFFER. See the doc comment on that method (marked "LOCAL PATCH +(breadcast, not upstream)"). ## Rolling the pin From 99a04c800bfe0af47146f29b0f7a962c41d5c095 Mon Sep 17 00:00:00 2001 From: Breadway Date: Sun, 16 Aug 2026 14:19:42 +0800 Subject: [PATCH 2/2] Route vX.Y.Z-rc.N tags to the beta track, not stable release.yml matches v*, which would have published an RC into the stable bakery index and pointed latest at it. Skip any tag whose name contains -rc there, and fire beta-release.yml on v*-rc.* instead. --- .forgejo/workflows/beta-release.yml | 26 ++++++++++++++------------ .forgejo/workflows/release.yml | 3 +++ AGENTS.md | 8 ++++---- 3 files changed, 21 insertions(+), 16 deletions(-) diff --git a/.forgejo/workflows/beta-release.yml b/.forgejo/workflows/beta-release.yml index 578cae6..08f8c7b 100644 --- a/.forgejo/workflows/beta-release.yml +++ b/.forgejo/workflows/beta-release.yml @@ -1,11 +1,11 @@ name: beta release -# Publishes a beta-track build on every push to `beta` — a frozen -# stabilization branch cut from `dev` when ready to stabilize; only -# fix/ branches merged into `beta` should land here afterward. -# See bread-ecosystem's docs/release-channels.md for the three-track policy. +# Publishes a beta-track build from a `vX.Y.Z-rc.N` tag. The leftover +# `beta` branch trigger is kept so an old branch push does not go silent; +# do not use that branch — cut beta from an RC tag. See CONTRIBUTING.md. on: push: + tags: ['v*-rc.*'] branches: ['beta'] jobs: @@ -16,7 +16,7 @@ jobs: run: | set -euo pipefail rm -rf src && mkdir src - git clone --branch beta --depth 1 \ + git clone --branch "${GITHUB_REF_NAME}" --depth 1 \ "https://git.breadway.dev/${GITHUB_REPOSITORY}.git" src - name: build @@ -25,19 +25,21 @@ jobs: - name: compute beta version run: | set -euo pipefail + # An RC tag is already valid semver and is the beta version. + # The leftover `beta` branch path still synthesizes one from the + # latest stable tag so an old push does not publish as 0.0.0. + if [[ "${GITHUB_REF_NAME}" == v*-rc.* ]]; then + echo "VERSION=${GITHUB_REF_NAME#v}" >> "$GITHUB_ENV" + exit 0 + fi cd src - # Base the beta version off the latest published stable tag, - # not Cargo.toml — Cargo.toml can go stale relative to the last - # real release (seen in practice: breadbox/breadpad/breadcrumbs/ - # breadpaper), which would make a beta build sort as OLDER than - # what's already installed and bakery would correctly refuse it. LATEST_TAG="$(git ls-remote --tags --refs \ "https://git.breadway.dev/${GITHUB_REPOSITORY}.git" 'v*' \ - | awk -F/ '{print $NF}' | sed 's/^v//' | sort -V | tail -1)" + | awk -F/ '{print $NF}' | sed 's/^v//' | grep -v -- '-rc' | sort -V | tail -1)" if [ -n "${LATEST_TAG}" ]; then CUR="${LATEST_TAG}" else - CUR="$(grep -m1 '^version' breadcast/Cargo.toml | sed -E 's/.*"(.*)".*/\1/')" + CUR="$(grep -m1 '^version' Cargo.toml | sed -E 's/.*"(.*)".*/\1/')" fi IFS='.' read -r MA MI PA <<< "${CUR}" SHA="$(git rev-parse --short HEAD)" diff --git a/.forgejo/workflows/release.yml b/.forgejo/workflows/release.yml index 9d08de9..bc62a22 100644 --- a/.forgejo/workflows/release.yml +++ b/.forgejo/workflows/release.yml @@ -6,6 +6,9 @@ on: jobs: build: + # RC tags are `v*` too (`v0.1.1-rc.1`) — those belong on the beta + # track, see beta-release.yml. A stable cut is a plain `vX.Y.Z`. + if: ${{ !contains(github.ref_name, '-rc') }} runs-on: [self-hosted, hestia] steps: - name: checkout diff --git a/AGENTS.md b/AGENTS.md index 0a3ec7e..f6c0ed9 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -20,12 +20,12 @@ out of sync with `dev`/`beta` across most repos in this ecosystem. - `github` — GitHub mirror. Push both when publishing. ## CI -- `release.yml` triggers on `push: tags: ['v*']` — tag a release to cut - the signed stable build. +- `release.yml` triggers on `push: tags: ['v*']` but skips any tag whose + name contains `-rc` — a plain `vX.Y.Z` cuts the signed stable build. +- `beta-release.yml` publishes a beta-track build from a `vX.Y.Z-rc.N` + tag (and still on leftover `beta` so an old branch push does not go silent). - `dev-release.yml` publishes a dev-track build on push to `main` (and still on leftover `dev` so an old branch push does not go silent). -- `beta-release.yml` still triggers on leftover `push: branches: ['beta']`. - Do **not** use that branch; cut beta from a `vX.Y.Z-rc.N` tag instead. - `check.yml` runs clippy + test on `main`, `feature/**`, and `fix/**`. ## Product cut