diff --git a/.forgejo/workflows/beta-release.yml b/.forgejo/workflows/beta-release.yml index 08f8c7b..578cae6 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 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. +# 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. 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 "${GITHUB_REF_NAME}" --depth 1 \ + git clone --branch beta --depth 1 \ "https://git.breadway.dev/${GITHUB_REPOSITORY}.git" src - name: build @@ -25,21 +25,19 @@ 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//' | grep -v -- '-rc' | sort -V | tail -1)" + | awk -F/ '{print $NF}' | sed 's/^v//' | sort -V | tail -1)" if [ -n "${LATEST_TAG}" ]; then CUR="${LATEST_TAG}" else - CUR="$(grep -m1 '^version' Cargo.toml | sed -E 's/.*"(.*)".*/\1/')" + CUR="$(grep -m1 '^version' breadcast/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/check.yml b/.forgejo/workflows/check.yml index 65ae3e6..b547c34 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: ['main', 'feature/**', 'fix/**'] + branches: ['feature/**', 'fix/**'] jobs: check: diff --git a/.forgejo/workflows/dev-release.yml b/.forgejo/workflows/dev-release.yml index 0207d3d..3d36231 100644 --- a/.forgejo/workflows/dev-release.yml +++ b/.forgejo/workflows/dev-release.yml @@ -1,8 +1,9 @@ name: dev release -# 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. +# 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. on: push: branches: ['main', 'dev'] @@ -15,7 +16,7 @@ jobs: run: | set -euo pipefail rm -rf src && mkdir src - git clone --branch "${GITHUB_REF_NAME}" --depth 1 \ + git clone --branch main --depth 1 \ "https://git.breadway.dev/${GITHUB_REPOSITORY}.git" src - name: build @@ -36,7 +37,7 @@ jobs: if [ -n "${LATEST_TAG}" ]; then CUR="${LATEST_TAG}" else - CUR="$(grep -m1 '^version' Cargo.toml | sed -E 's/.*"(.*)".*/\1/')" + CUR="$(grep -m1 '^version' breadcast/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 bc62a22..9d08de9 100644 --- a/.forgejo/workflows/release.yml +++ b/.forgejo/workflows/release.yml @@ -6,9 +6,6 @@ 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 f6c0ed9..29ba49e 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -20,13 +20,15 @@ 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*']` 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). -- `check.yml` runs clippy + test on `main`, `feature/**`, and `fix/**`. +- `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`. ## Product cut breadcast is a **bakery product**, not shipped on the BOS ISO, and not part @@ -39,9 +41,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` 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"`. +`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"`. ## Don't - Don't embed credentials in remote URLs — SSH or a credential helper only. diff --git a/EVENTS.md b/EVENTS.md index 6b76d57..d198d1a 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 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.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.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 a6416b0..68e5c5d 100644 --- a/bakery.toml +++ b/bakery.toml @@ -1,9 +1,7 @@ name = "breadcast" description = "Cast your screen to any Chromecast/Google TV or DLNA renderer — daemon + GTK4 popup" binaries = ["breadcast", "breadcastd"] -# 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 +# 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 @@ -16,7 +14,6 @@ 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 4176a31..3ec4288 100644 --- a/breadcast-caststream-sys/src/facade.cc +++ b/breadcast-caststream-sys/src/facade.cc @@ -209,20 +209,9 @@ 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::PlatformClientPosix* instance = - openscreen::PlatformClientPosix::GetInstance(); - if (instance == nullptr) { - return nullptr; - } - openscreen::TaskRunner& task_runner = instance->GetTaskRunner(); + openscreen::TaskRunner& task_runner = + openscreen::PlatformClientPosix::GetInstance()->GetTaskRunner(); breadcast_caststream::VideoParams params; params.width = width; @@ -266,13 +255,6 @@ 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) { @@ -290,10 +272,6 @@ 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); @@ -311,10 +289,6 @@ 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; } @@ -404,12 +378,6 @@ 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); @@ -450,18 +418,7 @@ void breadcast_caststream_sender_take_stats(CastStreamSender* sender, } int32_t breadcast_caststream_sender_needs_key_frame(CastStreamSender* sender) { - 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; + return sender->needs_key_frame.load(std::memory_order_relaxed) ? 1 : 0; } int32_t breadcast_caststream_sender_estimated_bandwidth_bps(CastStreamSender* sender) { @@ -480,24 +437,18 @@ 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. 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(); - } + // 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(); - if (openscreen::PlatformClientPosix::GetInstance() != nullptr) { - openscreen::PlatformClientPosix::ShutDown(); - } + openscreen::PlatformClientPosix::ShutDown(); delete sender; } diff --git a/breadcast-caststream-sys/src/facade.h b/breadcast-caststream-sys/src/facade.h index 840d890..8411721 100644 --- a/breadcast-caststream-sys/src/facade.h +++ b/breadcast-caststream-sys/src/facade.h @@ -88,18 +88,15 @@ void breadcast_caststream_sender_on_message(CastStreamSender* sender, const char* message, size_t message_len); -// 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(). +// 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). 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 222bb6f..5837002 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-core::caststream` for the ergonomic, thread-safe -//! wrapper most callers should use instead. +//! unsafe -- see `breadcast-caststream` (not this crate) 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 4241c97..77b99cc 100644 --- a/breadcast-core/src/capture/portal.rs +++ b/breadcast-core/src/capture/portal.rs @@ -117,10 +117,6 @@ 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 9815cfd..433efa1 100644 --- a/breadcast-core/src/caststream.rs +++ b/breadcast-core/src/caststream.rs @@ -225,18 +225,16 @@ impl CastStreamSender { } } - /// 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. + /// 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. /// - /// 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. + /// 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). pub fn enqueue_frame(&self, data: &[u8], is_key_frame: bool, capture_time_us: i64) -> Result<()> { let result = unsafe { breadcast_caststream_sender_enqueue_frame( @@ -248,7 +246,7 @@ impl CastStreamSender { ) }; if result != 0 { - bail!("frame not posted (session not negotiated or shutting down)"); + bail!("frame not enqueued (session not negotiated yet)"); } Ok(()) } diff --git a/breadcast-core/src/discovery.rs b/breadcast-core/src/discovery.rs index 00a8495..35b8cda 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,10 +50,7 @@ fn pick_address(addresses: &[(ScopedIp, Instant)]) -> Option { }) }) .or_else(|| candidates.first()) - // 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()) + .map(|(ip, _)| ip.to_ip_addr()) } const SERVICE_TYPE: &str = "_googlecast._tcp.local."; @@ -143,21 +140,10 @@ impl Discovery { } } - let Some(host) = pick_address(addrs) else { + let Some(host) = pick_address(addrs).map(|ip| ip.to_string()) 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 79c1fef..9e44a07 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, MEDIA_RENDERER}; +use super::AV_TRANSPORT; /// 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(MEDIA_RENDERER); + let search_target = SearchTarget::URN(AV_TRANSPORT); match rupnp::discover(&search_target, SEARCH_TIMEOUT, None).await { Ok(stream) => { use futures_util::StreamExt; @@ -86,10 +86,6 @@ 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 e6d78f8..26b2199 100644 --- a/breadcast-core/src/dlna/mod.rs +++ b/breadcast-core/src/dlna/mod.rs @@ -21,7 +21,3 @@ 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 31b0e1d..2447ce3 100644 --- a/breadcast-core/src/dlna/session.rs +++ b/breadcast-core/src/dlna/session.rs @@ -56,36 +56,18 @@ 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}{metadata}" + "0{escaped}" ); self.service .action(&self.device_url, "SetAVTransportURI", &set_uri_payload) .await .context("SetAVTransportURI failed")?; - // Some renderers auto-play after SetURI and then reject Play; - // others stay TRANSITIONING for a moment. Either PLAYING state is - // success. - match self - .service + self.service .action(&self.device_url, "Play", "01") .await - { - Ok(_) => {} - Err(e) => { - if self.transport_state().await.ok().as_deref() != Some("PLAYING") { - return Err(e).context("Play failed"); - } - } - } + .context("Play failed")?; Ok(()) } @@ -124,7 +106,7 @@ impl DlnaSession { } } -pub(crate) fn xml_escape(input: &str) -> String { +fn xml_escape(input: &str) -> String { let mut escaped = String::with_capacity(input.len()); for c in input.chars() { match c { @@ -138,14 +120,3 @@ pub(crate) 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 6907c07..4565514 100644 --- a/breadcast-core/src/http_server.rs +++ b/breadcast-core/src/http_server.rs @@ -16,9 +16,8 @@ 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. [`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. +/// 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). /// /// Every servable path is namespaced under a random token /// (`//playlist.m3u8`, etc. — see [`HttpServer::token`]) rather than @@ -38,7 +37,6 @@ const WORKER_THREADS: usize = 8; pub struct HttpServer { addr: SocketAddr, token: String, - server: Option>, } impl HttpServer { @@ -76,19 +74,7 @@ impl HttpServer { }); } - 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(); - } + Ok(Self { addr, token }) } /// The bound address, e.g. `0.0.0.0:41823`. Combine with this @@ -104,17 +90,7 @@ 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 { - 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(); + format!("http://{host}:{}/{}/{relative}", self.addr.port(), self.token) } } @@ -147,7 +123,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). -pub(crate) fn parse_range(value: &str, len: usize) -> Option<(usize, usize)> { +fn parse_range(value: &str, len: usize) -> Option<(usize, usize)> { let spec = value.strip_prefix("bytes=")?; if spec.contains(',') || len == 0 { return None; @@ -257,11 +233,7 @@ fn handle_request(request: tiny_http::Request, root: &Path, token: &str) -> Resu ]; let (status, body) = match range { - // 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); + Some((start, end)) if start <= end && end < data.len() => { headers.push(("Content-Range".to_string(), format!("bytes {start}-{end}/{}", data.len()))); (206u16, data[start..=end].to_vec()) } @@ -292,25 +264,3 @@ 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 4c74838..fe20e75 100644 --- a/breadcast-core/src/net.rs +++ b/breadcast-core/src/net.rs @@ -2,8 +2,6 @@ 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 @@ -20,42 +18,8 @@ const EXCLUDED_PREFIXES: &[&str] = &["tailscale", "wg", "docker", "veth", "br-", /// 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 { - pick_lan_ip(None) -} + const EXCLUDED_PREFIXES: &[&str] = &["tailscale", "wg", "docker", "veth", "br-", "virbr", "lo"]; -/// 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() @@ -65,7 +29,7 @@ fn lan_ipv4_candidates() -> Result> { } let text = String::from_utf8_lossy(&output.stdout); - let mut candidates: Vec<(String, Ipv4Addr, u8)> = Vec::new(); + let mut candidates: Vec<(String, Ipv4Addr)> = 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(); @@ -78,41 +42,20 @@ fn lan_ipv4_candidates() -> Result> { continue; } let Some(cidr) = fields.next() 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); + let Some(addr) = cidr.split('/').next().and_then(|a| a.parse::().ok()) else { continue }; if !addr.is_private() { continue; } - candidates.push((iface.to_string(), addr, prefix.min(32))); + candidates.push((iface.to_string(), addr)); } - 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) -} + // 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"))); -#[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)); - } + 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 ced7f4e..fe46558 100644 --- a/breadcast-core/src/pipeline/mod.rs +++ b/breadcast-core/src/pipeline/mod.rs @@ -451,7 +451,6 @@ 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. @@ -459,18 +458,16 @@ pub fn pull_encoded_frame(appsink: &gst_app::AppSink) -> Result, return Ok(None); } // Teardown from another thread (`CastMirrorSession::stop` sets - // 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) - { + // 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) { return Ok(None); } - stalled_for += pulled_at.elapsed(); + stalled_for += CAPTURE_STALL_POLL; if stalled_for < CAPTURE_STALL_TIMEOUT { continue; } @@ -532,56 +529,11 @@ 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 daemon uses [`run_until_eos_or_error`]. +/// synchronous `main`/example; the real daemon will want an async/watch-based +/// version instead of blocking a thread. 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 73b7cb9..0f94064 100644 --- a/breadcast/src/css.rs +++ b/breadcast/src/css.rs @@ -47,11 +47,6 @@ 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 64250bb..3b1ccf3 100644 --- a/breadcast/src/ipc_client.rs +++ b/breadcast/src/ipc_client.rs @@ -29,8 +29,6 @@ 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 56202d6..85e9e02 100644 --- a/breadcast/src/main.rs +++ b/breadcast/src/main.rs @@ -13,9 +13,7 @@ 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, EventControllerKey, Label, ListBox, Orientation, SelectionMode, glib, -}; +use gtk4::{Align, Application, Box as GBox, Button, Label, ListBox, Orientation, SelectionMode, glib}; use ipc_client::IpcClient; const PANEL_WIDTH: i32 = 360; @@ -72,50 +70,12 @@ 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)); @@ -124,16 +84,10 @@ 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(); } }); @@ -146,26 +100,11 @@ 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 || { - 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; - } - } + while let Ok(message) = rx.try_recv() { + handle_server_message(message, &list, &empty_label, &status_pill, &stop_button); } glib::ControlFlow::Continue }); @@ -186,7 +125,6 @@ fn handle_server_message( message: ServerMessage, list: &ListBox, empty_label: &Label, - error_label: &Label, status_pill: &Label, stop_button: &Button, ) { @@ -200,8 +138,6 @@ 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, @@ -209,10 +145,7 @@ 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" => { - error_label.set_visible(false); - update_state(status_pill, stop_button, data); - } + "state_changed" => update_state(status_pill, stop_button, data), _ => {} } } diff --git a/breadcastd/src/cast_mirror.rs b/breadcastd/src/cast_mirror.rs index b8444ad..2990077 100644 --- a/breadcastd/src/cast_mirror.rs +++ b/breadcastd/src/cast_mirror.rs @@ -83,11 +83,7 @@ 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, - generation: u64, - ) -> Result { + pub async fn start(device: CastDevice, daemon_tx: tokio::sync::mpsc::Sender) -> Result { let capture = CaptureSession::start().await.context("failed to start portal screen capture")?; let video_node_id = capture.video_node_id(); @@ -98,18 +94,13 @@ 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) = - 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"); - } - }; + build_video_pipeline_for_streaming(video_node_id).context("failed to build the encode pipeline")?; { let pipeline_watch = pipeline.clone(); std::thread::spawn(move || { - match breadcast_core::pipeline::run_until_eos_or_error(&pipeline_watch) { + match breadcast_core::pipeline::run_until_error_or_timeout(&pipeline_watch, gst::ClockTime::from_seconds(3600)) + { Ok(outcome) => tracing::debug!(?outcome, "encode pipeline bus watcher ended"), Err(e) => tracing::error!(error = ?e, "encode pipeline error"), } @@ -120,37 +111,19 @@ 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 connect = tokio::task::spawn_blocking(move || { + let (session, _media_events, raw_messages) = tokio::task::spawn_blocking(move || { CastSession::connect_app( &device_for_connect, CastDeviceApp::Custom(breadcast_core::caststream::MIRRORING_APP_ID.to_string()), ) }) - .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"); - } - }; + .await + .context("connect_app task panicked")? + .context("failed to connect and launch the Mirroring receiver")?; let (sender, stream_events) = - 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"); - } - }; + CastStreamSender::start(&device.host, "sender-0", session.transport_id(), video_params) + .context("failed to start the Cast Streaming session")?; let sender = Arc::new(sender); let message_pump = { @@ -165,12 +138,9 @@ 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 { @@ -180,68 +150,37 @@ impl CastMirrorSession { } } CastStreamEvent::Negotiated => negotiated.store(true, Ordering::Release), - 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::Error(message) => tracing::warn!(%message, "Cast Streaming error"), 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"); - if let Some(sender) = started.sender.as_ref() { - sender.negotiate(); - } + sender.negotiate(); let deadline = tokio::time::Instant::now() + tokio::time::Duration::from_secs(10); - while !negotiated.load(Ordering::Acquire) - && !failed.load(Ordering::Acquire) - && tokio::time::Instant::now() < deadline - { + while !negotiated.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) { - started.stop().await; + // 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; anyhow::bail!("never received an ANSWER from {} (negotiation timed out)", device.name); } - if let Err(e) = started.pipeline.set_state(gst::State::Playing) { - started.stop().await; - return Err(e).context("failed to start the encode pipeline"); - } + pipeline.set_state(gst::State::Playing).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 = started.sender.as_ref().expect("sender installed above").clone(); + let sender = sender.clone(); std::thread::spawn(move || { let result = frame_pump_loop(&appsink, &encoder, &sender); if let Err(e) = result { @@ -251,12 +190,19 @@ 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 { generation }); + let _ = daemon_tx.blocking_send(DaemonCommand::SessionEnded); }) }; - started.frame_pump = Some(frame_pump); - Ok(started) + 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), + }) } /// Tears down the session. Order matters and is not interchangeable: diff --git a/breadcastd/src/daemon.rs b/breadcastd/src/daemon.rs index 1a164bb..95a0d30 100644 --- a/breadcastd/src/daemon.rs +++ b/breadcastd/src/daemon.rs @@ -33,29 +33,7 @@ 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 { 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, + SessionEnded, } pub fn spawn(events_tx: broadcast::Sender, bread_client: BreadClient) -> mpsc::Sender { @@ -67,8 +45,6 @@ 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, @@ -83,7 +59,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. -pub(crate) enum ActiveSession { +enum ActiveSession { Cast(CastMirrorSession), Dlna(Box), } @@ -110,13 +86,6 @@ 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 @@ -142,12 +111,6 @@ 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); @@ -156,32 +119,20 @@ impl Daemon { self.broadcast_state(); let _ = reply.send(Ok(())); } - 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() { + 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() { 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 @@ -196,8 +147,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); @@ -215,100 +166,54 @@ impl Daemon { } async fn start_cast(&mut self, device_id: String, reply: oneshot::Sender>) { - if self.active_session.is_some() || self.starting { + if self.active_session.is_some() { 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() { - 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, - }), + 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) => { // `{e:#}` (not `{e}`/`to_string()`) so the full anyhow // context chain reaches the caller/GUI instead of just // the outermost ".context()" message. - Err(e) => Err(StartFailed { device_id: device.id, error: format!("{e:#}") }), - }; - let _ = daemon_tx.send(DaemonCommand::StartFinished { generation, outcome, reply }).await; - }); + bread_events::emit_mirroring_failed(&self.bread_client, &device.id, &format!("{e:#}")); + let _ = reply.send(Err(format!("{e:#}"))); + } + } return; } if let Some(device) = self.dlna_devices.get(&device_id).cloned() { - 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, + 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(), protocol: Protocol::Dlna, - }), - Err(e) => Err(StartFailed { device_id: device.url, error: format!("{e:#}") }), - }; - let _ = daemon_tx.send(DaemonCommand::StartFinished { generation, outcome, reply }).await; - }); + }; + 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:#}"))); + } + } return; } - 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)); - } - } + let _ = reply.send(Err(format!("unknown device id \"{device_id}\""))); } fn device_list(&self) -> Vec { diff --git a/breadcastd/src/dlna_mirror.rs b/breadcastd/src/dlna_mirror.rs index 510cc00..c3c68e2 100644 --- a/breadcastd/src/dlna_mirror.rs +++ b/breadcastd/src/dlna_mirror.rs @@ -13,11 +13,8 @@ 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, run_until_eos_or_error, wait_for_playlist_segments, -}; +use breadcast_core::pipeline::{build_video_pipeline, hls_output_dir, wait_for_playlist_segments}; use breadcast_core::{CaptureSession, DlnaDevice, DlnaSession}; use gstreamer as gst; use gstreamer::prelude::*; @@ -30,20 +27,11 @@ 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 { @@ -55,57 +43,42 @@ 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, - generation: u64, - ) -> Result { + pub async fn start(device: DlnaDevice, daemon_tx: tokio::sync::mpsc::Sender) -> Result { let capture = CaptureSession::start().await.context("failed to start portal screen capture")?; let video_node_id = capture.video_node_id(); - 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"); - } - }; + 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")?; - 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"); + // 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"), + } + }); } - 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"); - } - }; + 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")?; // 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. 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"); - } - }; + // 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")?; let stream_url = http.url(lan_ip, "playlist.m3u8"); // Two segments, not three: this is a "don't hand the renderer a 404 @@ -114,42 +87,30 @@ 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. - 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"); - } + wait_for_playlist_segments(&output_dir.join("playlist.m3u8"), 2, Duration::from_secs(20)) + .await + .context("encode pipeline never produced playable HLS segments")?; - 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"); - } + 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")?; 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 { generation }).await; + let _ = daemon_tx.send(DaemonCommand::SessionEnded).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 { generation }).await; + let _ = daemon_tx.send(DaemonCommand::SessionEnded).await; return; } } @@ -157,69 +118,17 @@ impl DlnaMirrorSession { }) }; - 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, - }) + Ok(Self { pipeline, session, capture: Some(capture), poll_task }) } /// Tears down the session: stops polling, tells the renderer to stop, - /// stops the encode pipeline, shuts the HLS server, closes the portal - /// capture session, and deletes the recording directory. + /// stops the encode pipeline, and closes the portal capture session. 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"); @@ -227,45 +136,10 @@ 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() { - 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) = capture.close().await { + tracing::warn!(error = ?e, "failed to cleanly close the portal capture session"); } } - 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 703d50c..d8cbd85 100644 --- a/breadcastd/src/ipc.rs +++ b/breadcastd/src/ipc.rs @@ -13,8 +13,6 @@ //! 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; @@ -33,13 +31,7 @@ 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() { - // 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) + std::fs::create_dir_all(parent) .with_context(|| format!("failed to create socket dir {}", parent.display()))?; } // A stale socket file from an unclean previous exit makes bind() fail @@ -53,18 +45,7 @@ pub async fn serve(daemon_tx: mpsc::Sender, events_tx: broadcast: tracing::info!(path = %socket_path.display(), "IPC socket listening"); loop { - 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 (stream, _addr) = listener.accept().await.context("failed to accept IPC connection")?; let daemon_tx = daemon_tx.clone(); let events_rx = events_tx.subscribe(); tokio::spawn(async move { @@ -125,14 +106,10 @@ async fn handle_connection( continue; } }; - // 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); - }); + let response = handle_request(request, &daemon_tx).await; + if writer_tx.send(response).is_err() { + break; + } } forward_task.abort(); diff --git a/contrib/breadcastd.service b/contrib/breadcastd.service index 69571af..c7778b5 100644 --- a/contrib/breadcastd.service +++ b/contrib/breadcastd.service @@ -7,7 +7,7 @@ PartOf=graphical-session.target [Service] Type=simple -ExecStart=%h/.local/bin/breadcastd +ExecStart=%h/.cargo/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 94308d9..50f4b8f 100644 --- a/vendor/rust_cast-0.21.0/PATCHES.md +++ b/vendor/rust_cast-0.21.0/PATCHES.md @@ -22,14 +22,11 @@ channel claims the namespace. ## The patch -`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)"). +`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)"). ## Rolling the pin