Merge s2-host-fault: the screenshare host-fault path (S2)

A pixelpass host that dies mid-share is now torn down instead of
staying advertised: stdout EOF is synthesized as a terminal fault,
routed back into the core on a dedicated channel behind the reliable
arm of the biased select, gated by the ActiveShare generation so a
reaped child's late EOF is dropped as stale, and handled by retiring
the share — presence ticket removal and ScreenShareStopped ahead of
the reap wait, the explanatory error after.

Reviewed by Gemini (three rounds: branch review, full-range merge
review, fix verification round). Its P2s — Join's early-exit ordering
hole, the reap-then-presence advertising window, and the missing
presence-side gate — are fixed and mutation-verified. 640 lib tests;
three live gates green on the desktop, including a two-process
observer gate that reads the sharer's presence from a second real
node.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
2026-07-31 14:53:02 -04:00
co-authored by Claude Fable 5
4 changed files with 836 additions and 44 deletions
+12
View File
@@ -137,6 +137,16 @@ pub enum CoreCommand {
/// Stop sharing our screen: kill the pixelpass host and clear the presence
/// ticket. No-op when not sharing.
StopScreenShare,
/// **Core-internal.** The running pixelpass host's stdout ended — the
/// process died (or its event stream broke), so the share identified by
/// `generation` is over: reap the child, pull the ticket off presence, and
/// tell the user. Synthesized by the core's own notice-forwarder task; the
/// UI never sends it. `generation` scopes the fault to one specific host
/// spawn, so a stale fault (the user already stopped, or started a new
/// share) is ignored rather than tearing down the wrong share.
ScreenShareHostFault {
generation: u64,
},
/// Watch a peer's screen share: spawn a pixelpass viewer for `ticket` and
/// open it in a local player.
ViewShare {
@@ -271,6 +281,7 @@ pub fn delivery_class(cmd: &CoreCommand) -> DeliveryClass {
quality: _,
}
| CoreCommand::StopScreenShare
| CoreCommand::ScreenShareHostFault { generation: _ }
| CoreCommand::ViewShare {
ticket: _,
settings: _,
@@ -363,6 +374,7 @@ pub fn coalesce_key(cmd: &CoreCommand) -> Option<CoalesceKey> {
quality: _,
}
| CoreCommand::StopScreenShare
| CoreCommand::ScreenShareHostFault { generation: _ }
| CoreCommand::ViewShare {
ticket: _,
settings: _,
+144 -36
View File
@@ -1397,10 +1397,25 @@ async fn run_core_loop(
// later opt-in can immediately publish whatever is currently running.
let mut current_game: Option<crate::game::DetectedGame> = None;
let mut network_mode = NetworkMode::default();
// Pixelpass binary override (config), and the ticket of our own active screen
// share (rides our presence so the room — incl. late joiners — can watch).
// Pixelpass binary override (config), and our own active screen share: the
// ticket rides our presence so the room — incl. late joiners — can watch,
// and the generation ties host-fault notices to this specific host spawn
// (see `ScreenShareHostFault`). One variable on purpose: the ticket and the
// generation must appear and vanish together, or a stale fault could tear
// down a share it doesn't belong to.
let mut pixelpass_override: Option<String> = None;
let mut current_sharing: Option<String> = None;
struct ActiveShare {
generation: u64,
ticket: String,
}
let mut current_sharing: Option<ActiveShare> = None;
// Monotonic per-spawn counter feeding `ActiveShare::generation`.
let mut share_generations: u64 = 0;
// Host faults re-enter the loop here (the notice-forwarder task can't touch
// loop state). The loop keeps `host_fault_tx` to clone into each share's
// forwarder, so this channel never closes — the select arm's `Some` pattern
// is total in practice and a closed-channel branch would be unreachable.
let (host_fault_tx, mut host_fault_rx) = mpsc::unbounded_channel::<u64>();
let mut active_session: Option<ActiveSession> = None;
// Standalone capture-only mic meter, live only when no session exists.
@@ -1549,6 +1564,13 @@ async fn run_core_loop(
// reachable it is already covered — nothing to add here.
None => break,
},
// A share's notice-forwarder task reported the host's stdout ended.
// The `Some` pattern is total: this loop owns `host_fault_tx` (see
// its declaration), so the channel cannot close — no `None` arm is
// written because one would be unreachable by construction.
Some(generation) = host_fault_rx.recv() => {
CoreCommand::ScreenShareHostFault { generation }
}
game_change = next_game_change(&mut game_rx) => {
// The detector worker published a new debounced game (or `None`).
let Some(detected) = game_change else {
@@ -1567,7 +1589,7 @@ async fn run_core_loop(
let self_state = presence.to_state(
is_muted.load(Ordering::Relaxed),
net.endpoint.addr(),
current_sharing.clone(),
current_sharing.as_ref().map(|s| s.ticket.clone()),
);
let _ = session.room_state.update_self_state(self_state).await;
}
@@ -1711,6 +1733,13 @@ async fn run_core_loop(
net.file_router.clear();
*current_room.lock().unwrap() = None;
}
// Any advertised share died with that session — deliberately —
// so retire it HERE, before the invalid-ticket early exit below
// can skip it. Left populated, the killed host's stdout EOF
// would pass the ScreenShareHostFault staleness gate and
// surface as a spurious "ended unexpectedly" error on top of
// the ticket error (Gemini review of S2, P2-1).
current_sharing = None;
// If a network-mode / identity change was deferred while a call was
// active, rebuild the persistent stack now — after the old session is
@@ -1806,8 +1835,8 @@ async fn run_core_loop(
secret_key.clone(),
));
// Fresh join starts not sharing; clear any stale share ticket.
current_sharing = None;
// (The share was already retired beside the session teardown
// above; a fresh join starts not sharing.)
let self_state =
presence.to_state(is_muted.load(Ordering::Relaxed), endpoint.addr(), None);
@@ -2828,8 +2857,11 @@ async fn run_core_loop(
is_muted.store(new_state, Ordering::Relaxed);
if let Some(session) = &active_session {
let self_state =
presence.to_state(new_state, net.endpoint.addr(), current_sharing.clone());
let self_state = presence.to_state(
new_state,
net.endpoint.addr(),
current_sharing.as_ref().map(|s| s.ticket.clone()),
);
let _ = session.room_state.update_self_state(self_state).await;
}
}
@@ -2842,7 +2874,7 @@ async fn run_core_loop(
let self_state = presence.to_state(
is_muted.load(Ordering::Relaxed),
net.endpoint.addr(),
current_sharing.clone(),
current_sharing.as_ref().map(|s| s.ticket.clone()),
);
let _ = session.room_state.update_self_state(self_state).await;
}
@@ -3155,7 +3187,7 @@ async fn run_core_loop(
let self_state = presence.to_state(
is_muted.load(Ordering::Relaxed),
net.endpoint.addr(),
current_sharing.clone(),
current_sharing.as_ref().map(|s| s.ticket.clone()),
);
let _ = session.room_state.update_self_state(self_state).await;
}
@@ -3361,7 +3393,7 @@ async fn run_core_loop(
let self_state = presence.to_state(
is_muted.load(Ordering::Relaxed),
net.endpoint.addr(),
current_sharing.clone(),
current_sharing.as_ref().map(|s| s.ticket.clone()),
);
let _ = session.room_state.update_self_state(self_state).await;
}
@@ -3435,46 +3467,62 @@ async fn run_core_loop(
continue;
}
};
// Forward pixelpass `app_audio` events (only emitted when an app
// is selected) to the UI so it can warn when the chosen app's
// audio drops. The channel closes when the host dies (drain hits
// EOF), ending the forwarder task on its own.
let notices = audio_app.as_deref().map(|_| {
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<
crate::screenshare::PixelpassEvent,
>();
let ui_tx_notices = ui_tx.clone();
tokio::spawn(async move {
while let Some(ev) = rx.recv().await {
let active = match ev {
crate::screenshare::PixelpassEvent::AppAudioRouted => true,
crate::screenshare::PixelpassEvent::AppAudioLost => false,
_ => continue,
};
if ui_tx_notices
.send(UiEvent::ShareAudioActive(active))
.await
.is_err()
{
// Every share gets a notice forwarder — not just app-audio ones.
// pixelpass `app_audio` events (only emitted when an app is
// selected) become UI warnings, and the drain's terminal `Eof`
// becomes a host fault scoped to this spawn's generation, so a
// host that dies is torn down instead of staying advertised in
// presence forever. On a failed spawn the sender is dropped
// before the drain ever runs, so the forwarder just ends and no
// fault is sent (the spawn error carries the news instead).
share_generations += 1;
let generation = share_generations;
let (notices_tx, mut notices_rx) =
tokio::sync::mpsc::unbounded_channel::<crate::screenshare::HostNotice>();
let ui_tx_notices = ui_tx.clone();
let fault_tx = host_fault_tx.clone();
tokio::spawn(async move {
while let Some(notice) = notices_rx.recv().await {
match notice {
crate::screenshare::HostNotice::Event(ev) => {
let active = match ev {
crate::screenshare::PixelpassEvent::AppAudioRouted => true,
crate::screenshare::PixelpassEvent::AppAudioLost => false,
_ => continue,
};
if ui_tx_notices
.send(UiEvent::ShareAudioActive(active))
.await
.is_err()
{
break;
}
}
// Terminal by contract: nothing follows on the
// channel, so the task ends here.
crate::screenshare::HostNotice::Eof => {
let _ = fault_tx.send(generation);
break;
}
}
});
tx
}
});
match crate::screenshare::spawn_host(
&bin,
audio_app.as_deref(),
&settings,
quality,
notices,
notices_tx,
)
.await
{
Ok((child, ticket)) => {
crate::log_msg("Screen share host started");
session.teardown.set_host(child);
current_sharing = Some(ticket.clone());
current_sharing = Some(ActiveShare {
generation,
ticket: ticket.clone(),
});
let self_state = presence.to_state(
is_muted.load(Ordering::Relaxed),
net.endpoint.addr(),
@@ -3523,6 +3571,66 @@ async fn run_core_loop(
let _ = ui_tx.send(UiEvent::ScreenShareStopped).await;
}
CoreCommand::ScreenShareHostFault { generation } => {
// Stale unless it names the share we are advertising RIGHT NOW.
// Every deliberate end of a share (StopScreenShare, Leave, a
// fresh Join) clears `current_sharing` before or while reaping
// the child, and the reaped child's stdout EOF then arrives
// here late — dropping it is the correct handling, not an edge
// case. A mismatched generation likewise: that fault belongs to
// an older spawn than the share now running.
let stale = current_sharing.as_ref().map(|s| s.generation) != Some(generation);
if stale {
continue;
}
crate::log_msg(
"Screen share host died (stdout EOF with the share still advertised)",
);
current_sharing = None;
// Pull the ticket off presence FIRST, before the reap: if the
// child only closed stdout and lives on, `stop_host` burns the
// full stop grace before the SIGKILL fallback, and for that
// whole window peers would still see (and click Watch on) a
// share whose host is already gone (Gemini S2-merge review,
// P2-1).
if let Some(session) = &mut active_session {
let self_state = presence.to_state(
is_muted.load(Ordering::Relaxed),
net.endpoint.addr(),
None,
);
let _ = session.room_state.update_self_state(self_state).await;
}
// Stopped next — it clears the UI's sharing state — so the
// local UI also stops saying "sharing" before the reap wait,
// and the error explaining why comes only after, so the user
// is never left looking at a "sharing" UI with an error
// beside it.
let _ = ui_tx.send(UiEvent::ScreenShareStopped).await;
let mut unconfirmed = false;
if let Some(session) = &mut active_session {
// The child is usually already dead, so this confirms the
// reap immediately; if it merely closed stdout and lives
// on, this is the SIGINT → grace → SIGKILL path. Either
// way the dead-or-dying child leaves the teardown slot, so
// `is_sharing` stops lying.
unconfirmed = matches!(
session.teardown.stop_host().await,
Some(teardown::StopOutcome::Unconfirmed)
);
}
let detail = if unconfirmed {
" Its process also couldn't be confirmed dead — check for a stray pixelpass."
} else {
""
};
let _ = ui_tx
.send(UiEvent::Error(format!(
"Screen share ended unexpectedly — pixelpass exited.{detail}"
)))
.await;
}
CoreCommand::ViewShare { ticket, settings } => {
let bin = match crate::screenshare::pixelpass_path(pixelpass_override.as_deref()) {
Some(b) => b,
+91 -8
View File
@@ -90,6 +90,21 @@ pub enum PixelpassEvent {
Other,
}
/// What the host's stdout drain forwards to the core over the notice channel.
///
/// `Eof` is **synthesized here**, not parsed: pixelpass has no "I died" event,
/// and a crash can abort across `extern "C"` before any JSON line is written,
/// so the stream ending is the only reliable death signal. A read *error*
/// counts too — either way the event stream is gone and the host must be
/// treated as over.
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum HostNotice {
/// A parsed pixelpass event line.
Event(PixelpassEvent),
/// The host's stdout ended (EOF or read error). Terminal: nothing follows.
Eof,
}
/// Parse a single stdout line from pixelpass `--output json`. Pure: no I/O.
pub fn parse_pixelpass_event(line: &str) -> Option<PixelpassEvent> {
let line = line.trim();
@@ -370,7 +385,10 @@ pub fn is_available(config_override: Option<&str>) -> bool {
/// `audio_app` is `Some`, pixelpass captures only that app's audio instead of the
/// whole desktop sink, which avoids the call-loopback echo (A23). The child keeps
/// running (streaming to viewers) until killed or dropped; remaining stdout is
/// drained in a background task so a full pipe can't stall the host. We do
/// drained in a background task so a full pipe can't stall the host. The drain
/// forwards every parsed event over `notices` and — the part no share may opt
/// out of — a terminal [`HostNotice::Eof`] when the stream ends, which is the
/// caller's only reliable signal that the host died. We do
/// not pass encode/viewer overrides unless the local settings explicitly ask for
/// them, so pixelpass keeps its own defaults in the common case.
pub async fn spawn_host(
@@ -378,7 +396,7 @@ pub async fn spawn_host(
audio_app: Option<&str>,
settings: &ScreenShareSettings,
quality: ShareQuality,
notices: Option<tokio::sync::mpsc::UnboundedSender<PixelpassEvent>>,
notices: tokio::sync::mpsc::UnboundedSender<HostNotice>,
) -> std::io::Result<(Child, String)> {
let args = host_args(audio_app, settings, quality);
// Log the exact argv we hand pixelpass so a field log can confirm which
@@ -433,7 +451,7 @@ pub async fn spawn_host(
if let Some(stderr) = stderr {
drain_stderr_in_background(stderr);
}
drain_in_background(lines, "host", notices);
drain_in_background(lines, "host", Some(notices));
Ok((child, ticket))
}
@@ -572,13 +590,15 @@ where
/// Keep reading the child's stdout to EOF in the background so a full pipe can't
/// stall it; log notable events for diagnostics. When `notices` is `Some`, each
/// parsed event is also forwarded to the caller (the core, which translates the
/// `app_audio` ones into a UI warning); a send failure (receiver dropped) just
/// stops forwarding, draining continues. The task ends on EOF (child exited).
/// parsed event is also forwarded to the caller (the core), and when the stream
/// ends — EOF or read error, i.e. the child exited or its event stream broke —
/// a final [`HostNotice::Eof`] is sent so the caller learns the child is gone
/// (a host that dies must not stay advertised as sharing). A send failure
/// (receiver dropped) just stops forwarding, draining continues.
fn drain_in_background<R>(
mut lines: tokio::io::Lines<BufReader<R>>,
role: &'static str,
notices: Option<tokio::sync::mpsc::UnboundedSender<PixelpassEvent>>,
notices: Option<tokio::sync::mpsc::UnboundedSender<HostNotice>>,
) where
R: tokio::io::AsyncRead + Unpin + Send + 'static,
{
@@ -587,10 +607,14 @@ fn drain_in_background<R>(
if let Some(ev) = parse_pixelpass_event(&line) {
crate::log_msg(&format!("pixelpass {role}: {}", event_for_log(&ev)));
if let Some(tx) = &notices {
let _ = tx.send(ev);
let _ = tx.send(HostNotice::Event(ev));
}
}
}
if let Some(tx) = &notices {
crate::log_msg(&format!("pixelpass {role}: stdout ended"));
let _ = tx.send(HostNotice::Eof);
}
});
}
@@ -1376,4 +1400,63 @@ Install hint: sudo apt install gstreamer1.0-plugins-bad
#[cfg(not(windows))]
assert_eq!(candidates, vec![dir.join("pixelpass")]);
}
/// The host-fault contract, clean-exit half: events are forwarded in order
/// and the stream ending yields exactly one terminal [`HostNotice::Eof`],
/// after which the drain task drops its sender (the closed channel is what
/// ends the core's forwarder). A host that dies silently — EOF swallowed —
/// is the S2 defect: the dead share stays advertised in presence.
#[tokio::test]
async fn drain_forwards_events_then_synthesizes_eof_when_stdout_ends() {
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
let (read_half, mut write_half) = tokio::io::duplex(1024);
drain_in_background(BufReader::new(read_half).lines(), "test", Some(tx));
use tokio::io::AsyncWriteExt;
write_half
.write_all(b"{\"event\":\"app_audio\",\"state\":\"routed\"}\nnot json\n")
.await
.unwrap();
drop(write_half); // child exited: stdout EOF
assert_eq!(
rx.recv().await,
Some(HostNotice::Event(PixelpassEvent::AppAudioRouted))
);
// The non-JSON line is dropped, not forwarded.
assert_eq!(rx.recv().await, Some(HostNotice::Eof));
assert_eq!(rx.recv().await, None, "task ended and dropped the sender");
}
/// The host-fault contract, broken-stream half: a read *error* (not a tidy
/// EOF) must synthesize the same terminal `Eof` — the event stream is gone
/// either way, and only the drain task can tell the core so.
#[tokio::test]
async fn drain_synthesizes_eof_on_a_read_error_too() {
struct BrokenPipe;
impl tokio::io::AsyncRead for BrokenPipe {
fn poll_read(
self: std::pin::Pin<&mut Self>,
_cx: &mut std::task::Context<'_>,
_buf: &mut tokio::io::ReadBuf<'_>,
) -> std::task::Poll<std::io::Result<()>> {
std::task::Poll::Ready(Err(std::io::Error::other("stream broke")))
}
}
use tokio::io::AsyncReadExt;
// One good event line, then the stream breaks mid-read.
let reader =
std::io::Cursor::new(b"{\"event\":\"capture\",\"state\":\"started\"}\n".to_vec())
.chain(BrokenPipe);
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
drain_in_background(BufReader::new(reader).lines(), "test", Some(tx));
assert_eq!(
rx.recv().await,
Some(HostNotice::Event(PixelpassEvent::CaptureStarted))
);
assert_eq!(rx.recv().await, Some(HostNotice::Eof));
assert_eq!(rx.recv().await, None, "task ended and dropped the sender");
}
}
+589
View File
@@ -0,0 +1,589 @@
//! S2 exit gate: a pixelpass host that dies mid-share must be torn down —
//! reaped, pulled off presence, `ScreenShareStopped` emitted **before** the
//! explanatory error — and a host stopped *deliberately* must NOT produce that
//! error when its stdout EOF arrives late (the staleness gate).
//!
//! Drives the real core loop end to end through `CoreController`, with the
//! pixelpass override pointed at fake shell scripts: one that emits a ticket
//! and dies, one that emits a ticket and lives until signalled. This is the
//! only harness that reaches the core's fault handler — the command loop has
//! no unit seam — so these two halves are what kill the "forwarder drops the
//! Eof" and "handler ignores the generation" mutants.
//!
//! Live: joins a real (solo) room, so it needs a working audio backend and
//! network access for the endpoint bind.
//! `cargo test --test screenshare_host_fault -- --ignored`
#![cfg(unix)]
use std::os::unix::fs::PermissionsExt;
use std::path::PathBuf;
use std::time::Duration;
use peerspeak::core::CoreController;
use peerspeak::core::messages::{CoreCommand, UiEvent};
const EVENT_TIMEOUT: Duration = Duration::from_secs(20);
/// How long to listen for events that must NOT arrive. Comfortably past the
/// fake host's exit plus the drain/forwarder hop, so a stale fault that WOULD
/// be mishandled has arrived by the end of it.
const QUIET_WINDOW: Duration = Duration::from_secs(3);
/// Removes the fake-pixelpass dir even when an assertion panics mid-test
/// (a plain trailing `remove_dir_all` never runs on an unwind).
struct TempDir(PathBuf);
impl Drop for TempDir {
fn drop(&mut self) {
std::fs::remove_dir_all(&self.0).ok();
}
}
fn write_fake_pixelpass(dir: &std::path::Path, name: &str, body: &str) -> PathBuf {
let path = dir.join(name);
std::fs::write(&path, body).expect("write fake pixelpass");
std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o755))
.expect("chmod fake pixelpass");
path
}
/// Skip events until `pick` matches, panicking after [`EVENT_TIMEOUT`].
/// Unrelated events (identity, presence, chat plumbing) flow on this channel
/// too, so gates scan rather than assert exact sequences.
async fn wait_for<T>(
rx: &mut tokio::sync::mpsc::Receiver<UiEvent>,
what: &str,
mut pick: impl FnMut(&UiEvent) -> Option<T>,
) -> T {
let deadline = tokio::time::Instant::now() + EVENT_TIMEOUT;
loop {
let ev = tokio::time::timeout_at(deadline, rx.recv())
.await
.unwrap_or_else(|_| panic!("timed out waiting for {what}"))
.unwrap_or_else(|| panic!("ui channel closed waiting for {what}"));
if let Some(v) = pick(&ev) {
return v;
}
}
}
#[tokio::test]
#[ignore = "live: joins a real solo room (audio backend + network bind)"]
async fn a_dead_host_is_torn_down_and_a_clean_stop_stays_clean() {
let dir_guard =
TempDir(std::env::temp_dir().join(format!("peerspeak-hostfault-{}", std::process::id())));
let dir = dir_guard.0.clone();
std::fs::create_dir_all(&dir).unwrap();
// Half 1's host: emits its ticket, then dies on its own — the S2 defect
// scenario. Plain `sleep` (no exec) so the shell itself exits and closes
// stdout with no orphan holding the pipe.
let dying_host = write_fake_pixelpass(
&dir,
"pixelpass-dies",
"#!/bin/sh\necho '{\"event\":\"ticket\",\"value\":\"fake-ticket-dies\"}'\nsleep 1\n",
);
// Half 2's host: lives until signalled. `exec` so the SIGINT from Stop
// Share hits the sleep itself — the process dies AND its stdout closes,
// which is exactly what makes the late Eof arrive and exercise the
// staleness gate rather than vacuously never sending a fault.
let living_host = write_fake_pixelpass(
&dir,
"pixelpass-lives",
"#!/bin/sh\necho '{\"event\":\"ticket\",\"value\":\"fake-ticket-lives\"}'\nexec sleep 600\n",
);
let (ui_tx, mut ui_rx) = tokio::sync::mpsc::channel(256);
let controller = CoreController::new(ui_tx);
assert!(controller.send(CoreCommand::SetPixelpassPath(Some(
dying_host.to_string_lossy().into_owned()
))));
assert!(controller.send(CoreCommand::Join {
name: "host-fault-gate".into(),
ticket: "create".into(),
room_name: "s2".into(),
input_device: None,
output_device: None,
echo_cancellation: false,
avatar: Default::default(),
}));
wait_for(&mut ui_rx, "RoomJoined", |ev| match ev {
UiEvent::RoomJoined { .. } => Some(()),
UiEvent::Error(e) => panic!("join failed: {e}"),
_ => None,
})
.await;
// ── Half 1: the host dies mid-share ─────────────────────────────────────
assert!(controller.send(CoreCommand::StartScreenShare {
audio_app: None,
settings: Default::default(),
quality: Default::default(),
}));
wait_for(
&mut ui_rx,
"ScreenShareStarted (dying host)",
|ev| match ev {
UiEvent::ScreenShareStarted => Some(()),
UiEvent::Error(e) => panic!("share start failed: {e}"),
_ => None,
},
)
.await;
// The fake host exits ~1s in. The contract: ScreenShareStopped FIRST (it
// clears the UI's sharing state), the explanatory error only after.
wait_for(&mut ui_rx, "ScreenShareStopped after host death", |ev| {
match ev {
UiEvent::ScreenShareStopped => Some(()),
// An error arriving first is the exact ordering defect S2 fixes:
// the UI would show "sharing" next to the explanation.
UiEvent::Error(e) => panic!("error arrived before ScreenShareStopped: {e}"),
_ => None,
}
})
.await;
let err = wait_for(&mut ui_rx, "the host-death error", |ev| match ev {
UiEvent::Error(e) => Some(e.clone()),
_ => None,
})
.await;
assert!(
err.contains("unexpectedly"),
"the error should say the share ended unexpectedly, got: {err}"
);
// ── Half 2: a deliberate stop must stay clean ───────────────────────────
assert!(controller.send(CoreCommand::SetPixelpassPath(Some(
living_host.to_string_lossy().into_owned()
))));
assert!(controller.send(CoreCommand::StartScreenShare {
audio_app: None,
settings: Default::default(),
quality: Default::default(),
}));
wait_for(
&mut ui_rx,
"ScreenShareStarted (living host)",
|ev| match ev {
UiEvent::ScreenShareStarted => Some(()),
UiEvent::Error(e) => panic!("second share start failed: {e}"),
_ => None,
},
)
.await;
assert!(controller.send(CoreCommand::StopScreenShare));
wait_for(
&mut ui_rx,
"ScreenShareStopped after Stop Share",
|ev| match ev {
UiEvent::ScreenShareStopped => Some(()),
UiEvent::Error(e) => panic!("clean stop produced an error: {e}"),
_ => None,
},
)
.await;
// The stopped host's stdout EOF is arriving about now as a *stale* fault
// (its generation was retired when Stop Share cleared the share). Without
// the staleness gate the handler would emit a second ScreenShareStopped
// and a spurious "ended unexpectedly" error — listen long enough for that
// mishandling to have shown up, and require silence.
let deadline = tokio::time::Instant::now() + QUIET_WINDOW;
while let Ok(Some(ev)) = tokio::time::timeout_at(deadline, ui_rx.recv()).await {
match ev {
UiEvent::ScreenShareStopped => {
panic!("stale host fault re-emitted ScreenShareStopped after a clean stop")
}
UiEvent::Error(e) if e.contains("unexpectedly") => {
panic!("stale host fault surfaced as an error after a clean stop: {e}")
}
_ => {}
}
}
// ── Half 3: a failed room switch while sharing must not cry "crash" ─────
// Join tears the old session down (killing the host, deliberately) BEFORE
// it validates the ticket, so an invalid ticket exits the Join arm early.
// The share must be retired at the teardown itself — left advertised, the
// killed host's EOF passes the staleness gate and a spurious "ended
// unexpectedly" lands on top of the ticket error (Gemini review, P2-1).
assert!(controller.send(CoreCommand::StartScreenShare {
audio_app: None,
settings: Default::default(),
quality: Default::default(),
}));
wait_for(
&mut ui_rx,
"ScreenShareStarted (before failed switch)",
|ev| match ev {
UiEvent::ScreenShareStarted => Some(()),
UiEvent::Error(e) => panic!("third share start failed: {e}"),
_ => None,
},
)
.await;
assert!(controller.send(CoreCommand::Join {
name: "host-fault-gate".into(),
ticket: "definitely-not-a-ticket".into(),
room_name: "s2".into(),
input_device: None,
output_device: None,
echo_cancellation: false,
avatar: Default::default(),
}));
wait_for(&mut ui_rx, "the invalid-ticket error", |ev| match ev {
UiEvent::Error(e) if e.contains("invalid room ticket") => Some(()),
UiEvent::Error(e) => panic!("unexpected error before the ticket error: {e}"),
_ => None,
})
.await;
// The deliberately-killed host's EOF is arriving about now; it must be
// dropped as stale, not reported as a crash.
let deadline = tokio::time::Instant::now() + QUIET_WINDOW;
while let Ok(Some(ev)) = tokio::time::timeout_at(deadline, ui_rx.recv()).await {
match ev {
UiEvent::ScreenShareStopped => {
panic!("failed room switch re-emitted ScreenShareStopped for the torn-down share")
}
UiEvent::Error(e) if e.contains("unexpectedly") => {
panic!("deliberate teardown during a failed room switch reported as a crash: {e}")
}
_ => {}
}
}
}
/// S2 presence gate: a host fault must pull the share ticket off PRESENCE —
/// what remote peers actually see — and must do it BEFORE the reap wait, not
/// after. Nothing on the sharer's own `UiEvent` channel can witness either
/// half (presence is only observable from another node), so this test runs a
/// real second core as an OBSERVER and asserts the sharer's `PeerState.sharing`
/// goes `Some` → `None` on fault.
///
/// The observer runs in a SEPARATE PROCESS (`presence_probe_helper`, this same
/// test binary re-invoked): two in-process cores would load the same
/// `identity.key` and collapse into one node id, and swapping `XDG_CONFIG_HOME`
/// between spawns in-process races other threads' getenv.
///
/// The fake host is a WEDGE — it closes stdout (the fault) but ignores SIGINT
/// and lives until the SIGKILL fallback — so `stop_host` burns the full 2 s
/// grace and TIME becomes the discriminator, exactly like the SIGINT gate:
/// with presence-removal-first the observer sees the ticket clear ~1 s after
/// it appeared (the wedge's pre-fault lifetime); with the old
/// reap-then-presence ordering, only after ~3 s. The bound also makes the
/// "presence removal deleted" mutant fail by timeout instead of passing
/// vacuously.
///
/// Live: two real solo-room cores (audio backend + network bind each).
#[tokio::test]
#[ignore = "live: two real cores in one room (audio backend + network bind), observer subprocess"]
async fn a_host_fault_pulls_the_ticket_off_presence_within_the_grace() {
/// Mirrors `core::teardown::STOP_GRACE` (private): the wait the wedge
/// forces before the SIGKILL fallback reaps it.
const STOP_GRACE_MS: u128 = 2000;
let dir_guard = TempDir(
std::env::temp_dir().join(format!("peerspeak-presence-gate-{}", std::process::id())),
);
let dir = dir_guard.0.clone();
std::fs::create_dir_all(&dir).unwrap();
// Emits its ticket, shares for ~1 s, then closes stdout (the fault) while
// staying alive and ignoring SIGINT, so the reap must wait out the grace.
// The trailing sleep is NOT exec'd on purpose: it forks after stdout is
// closed, so it holds no pipe (the vacuous-staleness trap doesn't apply),
// and it merely idles out after the SIGKILL reaps the shell.
//
// The fake ticket must pass `screenshare::sanitize_ticket` (`endpoint` +
// alphanumerics): the OBSERVER's gossip ingest sanitizes peer-advertised
// tickets, and a garbage one is nulled to `sharing: None` there — the
// probe would never see the share appear and the gate would go vacuous.
let wedged_host = write_fake_pixelpass(
&dir,
"pixelpass-wedges",
"#!/bin/sh\ntrap '' INT\n\
echo '{\"event\":\"ticket\",\"value\":\"endpointaabwxjexzensznfvuudiapn5tyzws3angd2merarm\"}'\n\
sleep 1\nexec 1>&-\nsleep 30\n",
);
let (ui_tx, mut ui_rx) = tokio::sync::mpsc::channel(256);
let controller = CoreController::new(ui_tx);
assert!(controller.send(CoreCommand::SetPixelpassPath(Some(
wedged_host.to_string_lossy().into_owned()
))));
assert!(controller.send(CoreCommand::Join {
name: "presence-gate".into(),
ticket: "create".into(),
room_name: "s2-presence".into(),
input_device: None,
output_device: None,
echo_cancellation: false,
avatar: Default::default(),
}));
let room_ticket = wait_for(&mut ui_rx, "RoomJoined", |ev| match ev {
UiEvent::RoomJoined { ticket, .. } => Some(ticket.clone()),
UiEvent::Error(e) => panic!("join failed: {e}"),
_ => None,
})
.await;
// The observer, in its own process with its own config dir (fresh
// identity). It prints `PROBE …` lines this test parses.
let probe_config = dir.join("probe-config");
std::fs::create_dir_all(&probe_config).unwrap();
let probe = tokio::process::Command::new(std::env::current_exe().unwrap())
.kill_on_drop(true)
.args([
"presence_probe_helper",
"--exact",
"--ignored",
"--nocapture",
])
.env("PEERSPEAK_PROBE_TICKET", &room_ticket)
.env("XDG_CONFIG_HOME", &probe_config)
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::piped())
.spawn()
.expect("spawn the presence probe");
// Only share once the probe is in the room, so it witnesses the ticket
// APPEARING before the fault clears it (otherwise `Some` → `None` could
// both predate its join and the gate would go vacuous).
wait_for(&mut ui_rx, "the probe's PeerJoined", |ev| match ev {
UiEvent::PeerJoined { .. } => Some(()),
UiEvent::Error(e) => panic!("waiting for the probe: {e}"),
_ => None,
})
.await;
assert!(controller.send(CoreCommand::StartScreenShare {
audio_app: None,
settings: Default::default(),
quality: Default::default(),
}));
wait_for(
&mut ui_rx,
"ScreenShareStarted (wedged host)",
|ev| match ev {
UiEvent::ScreenShareStarted => Some(()),
UiEvent::Error(e) => panic!("share start failed: {e}"),
_ => None,
},
)
.await;
// Sharer-side contract, unchanged by the reorder: Stopped first, the
// explanatory error only after.
wait_for(
&mut ui_rx,
"ScreenShareStopped after the wedge faults",
|ev| match ev {
UiEvent::ScreenShareStopped => Some(()),
UiEvent::Error(e) => panic!("error arrived before ScreenShareStopped: {e}"),
_ => None,
},
)
.await;
let err = wait_for(&mut ui_rx, "the host-death error", |ev| match ev {
UiEvent::Error(e) => Some(e.clone()),
_ => None,
})
.await;
assert!(
err.contains("unexpectedly"),
"the error should say the share ended unexpectedly, got: {err}"
);
let out = tokio::time::timeout(Duration::from_secs(60), probe.wait_with_output())
.await
.expect("probe process outlived its budget")
.expect("probe process wait");
let stdout = String::from_utf8_lossy(&out.stdout);
let stderr = String::from_utf8_lossy(&out.stderr);
assert!(
out.status.success(),
"probe failed ({}).\nstdout:\n{stdout}\nstderr:\n{stderr}",
out.status
);
let cleared_ms: u128 = stdout
.lines()
.find_map(|l| l.strip_prefix("PROBE sharing-cleared "))
.unwrap_or_else(|| {
panic!("probe never saw the ticket clear from presence.\nstdout:\n{stdout}")
})
.trim()
.parse()
.expect("probe delta should be integer millis");
// Presence-removal-first: ~1000 ms (the wedge's pre-fault lifetime).
// Reap-then-presence: ~3000 ms (lifetime + the full stop grace). The
// grace itself splits them with ~1 s of jitter headroom on each side.
assert!(
cleared_ms < STOP_GRACE_MS,
"presence kept advertising the dead share for {cleared_ms} ms after it appeared — \
at or past the wedge lifetime + stop grace, i.e. the ticket was only removed \
AFTER the reap wait instead of before it"
);
assert!(controller.send(CoreCommand::Leave));
}
/// Observer half of `a_host_fault_pulls_the_ticket_off_presence_within_the_grace`,
/// run BY that test as a subprocess. Standalone (no `PEERSPEAK_PROBE_TICKET` in
/// the env — e.g. a plain `--ignored` sweep) it is a no-op pass.
#[tokio::test]
#[ignore = "helper: spawned by the presence gate as a subprocess; standalone it no-ops"]
async fn presence_probe_helper() {
let Ok(room_ticket) = std::env::var("PEERSPEAK_PROBE_TICKET") else {
return;
};
let (ui_tx, mut ui_rx) = tokio::sync::mpsc::channel(256);
let controller = CoreController::new(ui_tx);
assert!(controller.send(CoreCommand::Join {
name: "presence-probe".into(),
ticket: room_ticket,
room_name: String::new(),
input_device: None,
output_device: None,
echo_cancellation: false,
avatar: Default::default(),
}));
wait_for(&mut ui_rx, "RoomJoined (probe)", |ev| match ev {
UiEvent::RoomJoined { .. } => Some(()),
UiEvent::Error(e) => panic!("probe join failed: {e}"),
_ => None,
})
.await;
// Watch the sharer's presence: record when its `sharing` ticket appears,
// report the delta when it clears. Timings on both ends are local-loopback
// arrival times, so the parent's bound compares like with like.
let deadline = tokio::time::Instant::now() + Duration::from_secs(30);
let mut seen_at: Option<std::time::Instant> = None;
loop {
let ev = tokio::time::timeout_at(deadline, ui_rx.recv())
.await
.expect("probe timed out watching for the sharing transition")
.expect("probe ui channel closed");
let sharing = match &ev {
UiEvent::PeerJoined { state, .. } | UiEvent::PeerUpdated { state, .. } => {
state.sharing.is_some()
}
_ => continue,
};
match (&seen_at, sharing) {
(None, true) => {
seen_at = Some(std::time::Instant::now());
println!("PROBE sharing-seen");
}
(Some(t0), false) => {
println!("PROBE sharing-cleared {}", t0.elapsed().as_millis());
break;
}
_ => {}
}
}
assert!(controller.send(CoreCommand::Leave));
}
/// The long-owed Stop Share SIGINT gate (0c half (ii)), against the REAL
/// pixelpass binary: a Stop Share must end the host through the graceful
/// SIGINT path — child exits within [`STOP_GRACE`], no SIGKILL fallback, no
/// "couldn't confirm" warning — because SIGKILL would skip pixelpass's own
/// teardown (it unloads its capture sink on the way out in sink-owning modes).
///
/// The fallback is indistinguishable from success in the event stream (both
/// end in a confirmed reap), so the discriminator is TIME: the fallback path
/// first waits out the full 2 s grace, while a host honouring SIGINT exits in
/// milliseconds. The bound asserts the stop completed inside the grace.
///
/// Live: needs `pixelpass` on `$PATH` plus a real solo room (audio + network).
#[tokio::test]
#[ignore = "live: real pixelpass host + a real solo room (audio backend, network bind)"]
async fn stop_share_ends_the_real_host_via_sigint_within_the_grace() {
/// Mirrors `core::teardown::STOP_GRACE` (private): the graceful wait
/// before the SIGKILL fallback.
const STOP_GRACE: Duration = Duration::from_secs(2);
let (ui_tx, mut ui_rx) = tokio::sync::mpsc::channel(256);
let controller = CoreController::new(ui_tx);
// No override: resolve the real binary from $PATH.
assert!(controller.send(CoreCommand::SetPixelpassPath(None)));
assert!(controller.send(CoreCommand::Join {
name: "sigint-gate".into(),
ticket: "create".into(),
room_name: "s2".into(),
input_device: None,
output_device: None,
echo_cancellation: false,
avatar: Default::default(),
}));
wait_for(&mut ui_rx, "RoomJoined", |ev| match ev {
UiEvent::RoomJoined { .. } => Some(()),
UiEvent::Error(e) => panic!("join failed: {e}"),
_ => None,
})
.await;
// Whole-desktop share: no viewers ever connect, so the real host sits idle
// after its ticket (capture starts on first viewer) — exactly the state a
// Stop Share most often hits.
assert!(controller.send(CoreCommand::StartScreenShare {
audio_app: None,
settings: Default::default(),
quality: Default::default(),
}));
wait_for(
&mut ui_rx,
"ScreenShareStarted (real pixelpass)",
|ev| match ev {
UiEvent::ScreenShareStarted => Some(()),
UiEvent::Error(e) => panic!("real pixelpass host failed to start: {e}"),
_ => None,
},
)
.await;
let stop_started = std::time::Instant::now();
assert!(controller.send(CoreCommand::StopScreenShare));
wait_for(
&mut ui_rx,
"ScreenShareStopped (real pixelpass)",
|ev| match ev {
UiEvent::ScreenShareStopped => Some(()),
// An Unconfirmed reap surfaces exactly this way; it means the
// SIGINT AND the SIGKILL both failed to end the host.
UiEvent::Error(e) => panic!("stop of the real host was not clean: {e}"),
_ => None,
},
)
.await;
let elapsed = stop_started.elapsed();
assert!(
elapsed < STOP_GRACE,
"stop took {elapsed:?} — at or past the {STOP_GRACE:?} grace, i.e. the \
SIGKILL fallback fired instead of pixelpass honouring SIGINT"
);
// And the late stdout EOF from the SIGINTed host must stay silent (same
// staleness contract the fake-host half pins).
let deadline = tokio::time::Instant::now() + QUIET_WINDOW;
while let Ok(Some(ev)) = tokio::time::timeout_at(deadline, ui_rx.recv()).await {
match ev {
UiEvent::ScreenShareStopped => {
panic!("stale fault from the SIGINTed real host re-emitted ScreenShareStopped")
}
UiEvent::Error(e) if e.contains("unexpectedly") => {
panic!("stale fault from the SIGINTed real host surfaced as an error: {e}")
}
_ => {}
}
}
assert!(controller.send(CoreCommand::Leave));
}