diff --git a/README.md b/README.md
index d7e8a7a..e3de9e4 100644
--- a/README.md
+++ b/README.md
@@ -2,6 +2,12 @@
A **mobile companion** for [KDE Connect](https://kdeconnect.kde.org) / [Valent](https://valent.andyholmes.ca/) peers.
+
+
+
+
+Waydroid (x86_64) — egui + vidya shell, on-device KDE Connect peer
+
| Layer | Role |
|--------|------|
| **Gleam** (`global_app`, `dist_driver`) | Domain model + controller over Erlang distribution |
diff --git a/dist-ui/src/egui_ui.rs b/dist-ui/src/egui_ui.rs
index 71ef277..c9c2322 100644
--- a/dist-ui/src/egui_ui.rs
+++ b/dist-ui/src/egui_ui.rs
@@ -94,17 +94,17 @@ impl App for DistApp {
}
}
- egui::TopBottomPanel::top("header")
- .frame(th.header_frame())
- .show(ctx, |ui| {
- ui.horizontal(|ui| {
- ui.spacing_mut().item_spacing = Vec2::new(th.spacing.md, th.spacing.sm);
- vidya::title(ui, th, &self.title);
- ui.with_layout(egui::Layout::right_to_left(egui::Align::Center), |ui| {
- vidya::dim_label(ui, th, &self.status);
- });
+ // Vidya owns system status/nav safe areas on Android (edge-to-edge).
+ // Prefer top_header so title/status cannot sit under the clock widgets.
+ vidya::top_header(ctx, th, |ui| {
+ ui.horizontal(|ui| {
+ ui.spacing_mut().item_spacing = Vec2::new(th.spacing.md, th.spacing.sm);
+ vidya::title(ui, th, &self.title);
+ ui.with_layout(egui::Layout::right_to_left(egui::Align::Center), |ui| {
+ vidya::dim_label(ui, th, &self.status);
});
});
+ });
egui::CentralPanel::default()
.frame(th.page_frame())
@@ -170,6 +170,13 @@ fn show_engine_panel(
if !snap.activity.is_empty() {
ui.add_space(th.spacing.xs);
vidya::dim_label(ui, th, &format!("Activity: {}", snap.activity));
+ // Keep status toast in sync with async ping results.
+ if snap.activity.starts_with("Ping ")
+ || snap.activity.starts_with("Sending ping")
+ || snap.activity.contains("channel")
+ {
+ *kc_msg = snap.activity.clone();
+ }
}
ui.add_space(th.spacing.sm);
@@ -200,16 +207,53 @@ fn show_engine_panel(
}
}
- if !snap.pending_pairs.is_empty() {
+ // Incoming pair requests (and a linked-but-untrusted peer for one-tap Accept).
+ let mut accept_targets: Vec<(String, String)> = snap.pending_pairs.clone();
+ if let (Some(id), Some(name)) = (&snap.linked_peer_id, &snap.linked_peer_name) {
+ let trusted = snap.trusted.iter().any(|(tid, _)| tid == id);
+ if !trusted && !accept_targets.iter().any(|(tid, _)| tid == id) {
+ accept_targets.push((id.clone(), name.clone()));
+ }
+ }
+ // Live unpaired peers not already listed.
+ for id in &snap.live_peers {
+ let trusted = snap.trusted.iter().any(|(tid, _)| tid == id);
+ if trusted || accept_targets.iter().any(|(tid, _)| tid == id) {
+ continue;
+ }
+ let name = snap
+ .nearby
+ .iter()
+ .find(|(nid, _, _)| nid == id)
+ .map(|(_, n, _)| n.clone())
+ .or_else(|| snap.linked_peer_name.clone())
+ .unwrap_or_else(|| id[..8.min(id.len())].to_string());
+ accept_targets.push((id.clone(), name));
+ }
+
+ if !accept_targets.is_empty() {
vidya::title(ui, th, "Pair requests");
- for (id, name) in &snap.pending_pairs {
+ vidya::dim_label(
+ ui,
+ th,
+ "Valent sent a pair request — Accept to finish pairing.",
+ );
+ for (id, name) in &accept_targets {
ui.horizontal(|ui| {
vidya::body(ui, th, &format!("From {name}"));
- if ui.button("Accept").clicked() {
+ if vidya::primary_button(ui, th, "Accept pair").clicked() {
if let Some(eng) = engine::global_engine() {
*kc_msg = match eng.accept_pair(id) {
Ok(s) => s,
- Err(e) => format!("Failed: {e}"),
+ Err(e) => format!("Accept failed: {e}"),
+ };
+ }
+ }
+ if ui.button("Reject").clicked() {
+ if let Some(eng) = engine::global_engine() {
+ *kc_msg = match eng.unpair(id) {
+ Ok(s) => s,
+ Err(e) => format!("Reject failed: {e}"),
};
}
}
@@ -220,17 +264,22 @@ fn show_engine_panel(
vidya::title(ui, th, "Trusted peers");
if snap.trusted.is_empty() {
- vidya::dim_label(ui, th, "None yet — use Request pairing in Valent.");
+ vidya::dim_label(
+ ui,
+ th,
+ "None yet — pair from Valent, then Accept pair above if needed.",
+ );
} else {
for (id, name) in &snap.trusted {
let live = snap.live_peers.iter().any(|p| p == id);
ui.horizontal(|ui| {
- let mark = if live { "●" } else { "○" };
+ // Drawn circle — Unicode ●/○ becomes a box (missing glyph) on Android.
+ vidya::status_dot(ui, th, live);
vidya::body(
ui,
th,
&format!(
- "{mark} {name} ({}){}",
+ "{name} ({}){}",
&id[..8.min(id.len())],
if live { " · live" } else { " · offline" }
),
@@ -240,7 +289,9 @@ fn show_engine_panel(
ui.add_enabled_ui(live, |ui| {
if ui.button("Ping").clicked() {
if let Some(eng) = engine::global_engine() {
- *kc_msg = eng.ping_peer(id);
+ // Never block the UI/render thread on TLS (Android ANR).
+ *kc_msg = "Sending ping…".into();
+ eng.ping_peer_async(id.clone());
}
}
});
diff --git a/dist-ui/src/kdeconnect/channel.rs b/dist-ui/src/kdeconnect/channel.rs
index 54e189c..a030494 100644
--- a/dist-ui/src/kdeconnect/channel.rs
+++ b/dist-ui/src/kdeconnect/channel.rs
@@ -114,11 +114,12 @@ fn configure_stream(stream: &TcpStream) -> Result<(), String> {
stream
.set_nonblocking(false)
.map_err(|e| e.to_string())?;
- // No short read timeout: a 10s timeout made idle paired channels look
- // "Disconnected" in Valent (EAGAIN/TimedOut treated as hangup).
+ // Session loop sets a short read timeout while polling. Keep default None so
+ // idle channels are not closed as TimedOut (Valent "Disconnected").
stream.set_read_timeout(None).map_err(|e| e.to_string())?;
+ // Short write timeout: UI must never block for seconds (Android ANR).
stream
- .set_write_timeout(Some(Duration::from_secs(30)))
+ .set_write_timeout(Some(Duration::from_secs(2)))
.map_err(|e| e.to_string())?;
let _ = stream.set_nodelay(true);
Ok(())
@@ -362,13 +363,26 @@ pub fn write_packet(stream: &mut S, packet: &KdePacket) -> Result<(),
stream.flush().map_err(|e| e.to_string())
}
+/// Read one newline-terminated packet.
+///
+/// On idle sockets (read timeout / would-block with **no** bytes yet), returns
+/// `Err("TimedOut")` immediately so the caller can drain outbound work and retry.
+/// Never sleeps while holding the caller's stream lock.
pub fn read_packet_line(stream: &mut S) -> Result {
let mut buf = Vec::new();
let mut byte = [0u8; 1];
+ // Cap how long we sit mid-line waiting for more data (ms of consecutive timeouts).
+ let mut mid_line_timeouts: u32 = 0;
loop {
match stream.read(&mut byte) {
- Ok(0) => return Err("eof".into()),
+ Ok(0) => {
+ if buf.is_empty() {
+ return Err("eof".into());
+ }
+ return Err("eof mid-packet".into());
+ }
Ok(_) => {
+ mid_line_timeouts = 0;
if byte[0] == b'\n' {
break;
}
@@ -382,8 +396,17 @@ pub fn read_packet_line(stream: &mut S) -> Result {
|| e.kind() == std::io::ErrorKind::TimedOut
|| e.kind() == std::io::ErrorKind::Interrupted =>
{
- // Transient — keep the channel; caller may also retry.
- std::thread::sleep(Duration::from_millis(20));
+ if buf.is_empty() {
+ // Idle: let the session loop flush outbound pings.
+ return Err("TimedOut".into());
+ }
+ // Mid-line: allow a few short timeouts, then surface so outbound can drain.
+ mid_line_timeouts += 1;
+ if mid_line_timeouts > 40 {
+ // ~2s at 50ms socket timeout
+ return Err("TimedOut".into());
+ }
+ // Do not sleep here — socket timeout already waited.
continue;
}
Err(e) => return Err(e.to_string()),
diff --git a/dist-ui/src/kdeconnect/engine.rs b/dist-ui/src/kdeconnect/engine.rs
index 9adc007..6622f1a 100644
--- a/dist-ui/src/kdeconnect/engine.rs
+++ b/dist-ui/src/kdeconnect/engine.rs
@@ -17,13 +17,19 @@ use crate::kdeconnect::packet::KdePacket;
use crate::kdeconnect::pair::{PairEffect, PairEvent, PairState, PairStateMachine};
use crate::kdeconnect::trust::TrustStore;
-/// One live TLS session. Shared so UI can **write synchronously** while the
-/// session thread reads. Valent opens several connections; we keep all and
-/// fan-out sends.
+/// Work item for the session thread (never written from the UI thread).
+enum OutboundMsg {
+ Packet(KdePacket),
+ /// Packet + oneshot so the caller can wait for the TLS write result.
+ PacketAck(KdePacket, std::sync::mpsc::Sender>),
+}
+
+/// One live TLS session. The session thread owns reads/writes; UI only enqueues.
#[derive(Clone)]
struct SessionEntry {
gen: u64,
stream: Arc>,
+ outbound: std::sync::mpsc::Sender,
}
#[derive(Debug, Clone)]
@@ -41,6 +47,8 @@ pub struct EngineSnapshot {
/// Valent-style verification key (8 hex chars) while a channel is live.
pub verification_key: Option,
pub linked_peer_name: Option,
+ /// Device id of the currently linked peer (if known), for Accept UI.
+ pub linked_peer_id: Option,
pub status: String,
/// Last activity line for the UI (ping received, etc.).
pub activity: String,
@@ -60,6 +68,7 @@ pub struct Engine {
activity: Arc>,
verification_key: Arc>>,
linked_peer_name: Arc>>,
+ linked_peer_id: Arc>>,
stop: Arc,
packet_ids: Arc,
live_sessions: Arc,
@@ -115,6 +124,7 @@ impl Engine {
let live_sessions = Arc::new(AtomicU64::new(0));
let verification_key = Arc::new(Mutex::new(None));
let linked_peer_name = Arc::new(Mutex::new(None));
+ let linked_peer_id = Arc::new(Mutex::new(None));
let sessions: Arc>>> =
Arc::new(Mutex::new(HashMap::new()));
let session_gen = Arc::new(AtomicU64::new(1));
@@ -130,6 +140,7 @@ impl Engine {
let live_a = live_sessions.clone();
let vkey_a = verification_key.clone();
let lname_a = linked_peer_name.clone();
+ let lid_a = linked_peer_id.clone();
let sessions_a = sessions.clone();
let session_gen_a = session_gen.clone();
let activity_a = activity.clone();
@@ -145,6 +156,7 @@ impl Engine {
live_a,
vkey_a,
lname_a,
+ lid_a,
sessions_a,
session_gen_a,
activity_a,
@@ -164,6 +176,7 @@ impl Engine {
activity,
verification_key,
linked_peer_name,
+ linked_peer_id,
stop,
packet_ids,
live_sessions,
@@ -220,6 +233,11 @@ impl Engine {
.lock()
.ok()
.and_then(|g| g.clone());
+ let linked_peer_id = self
+ .linked_peer_id
+ .lock()
+ .ok()
+ .and_then(|g| g.clone());
let live_peers = self
.sessions
.lock()
@@ -241,6 +259,7 @@ impl Engine {
live_peers,
verification_key,
linked_peer_name,
+ linked_peer_id,
status,
activity,
}
@@ -259,66 +278,92 @@ impl Engine {
.channel_addr()
}
- /// Send a packet on **every** open primary channel for this peer (sync write).
- pub fn send_to_peer(&self, peer_id: &str, packet: KdePacket) -> Result {
- let streams: Vec>> = self
- .sessions
- .lock()
- .map_err(|e| e.to_string())?
- .get(peer_id)
- .map(|v| v.iter().map(|e| e.stream.clone()).collect())
- .unwrap_or_default();
- if streams.is_empty() {
- return Err(format!(
- "no live channel to {peer_id} — wait until Valent is connected"
- ));
- }
- let mut ok = 0usize;
- let mut last_err = String::new();
- for stream in streams {
- match stream.lock() {
- Ok(mut g) => match write_packet(&mut *g, &packet) {
- Ok(()) => {
- ok += 1;
- if let Ok(raw) = packet.encode_line() {
- eprintln!(
- "[kdeconnect] → {} (sync): {}",
- packet.packet_type,
- String::from_utf8_lossy(&raw).trim()
- );
- }
+ /// Enqueue a packet on the **newest** live channel for this peer.
+ ///
+ /// The session thread drains `outbound` and performs the TLS write.
+ /// If `wait` is true, block (with timeout) until the write completes.
+ pub fn send_to_peer(
+ &self,
+ peer_id: &str,
+ packet: KdePacket,
+ wait: bool,
+ ) -> Result {
+ let tx = {
+ let map = self.sessions.lock().map_err(|e| e.to_string())?;
+ let list = map.get(peer_id).ok_or_else(|| {
+ format!("no live channel to {peer_id} — wait until Valent is connected")
+ })?;
+ // Prefer newest generation (Valent thrash leaves dead older sockets).
+ list.iter()
+ .max_by_key(|e| e.gen)
+ .map(|e| e.outbound.clone())
+ .ok_or_else(|| format!("no live channel to {peer_id}"))?
+ };
+
+ if wait {
+ let (ack_tx, ack_rx) = std::sync::mpsc::channel();
+ tx.send(OutboundMsg::PacketAck(packet.clone(), ack_tx))
+ .map_err(|e| format!("session gone: {e}"))?;
+ match ack_rx.recv_timeout(Duration::from_secs(3)) {
+ Ok(Ok(())) => {
+ if let Ok(raw) = packet.encode_line() {
+ eprintln!(
+ "[kdeconnect] → {} (on wire): {}",
+ packet.packet_type,
+ String::from_utf8_lossy(&raw).trim()
+ );
}
- Err(e) => last_err = e,
- },
- Err(e) => last_err = e.to_string(),
+ Ok(1)
+ }
+ Ok(Err(e)) => Err(e),
+ Err(_) => Err("ping timed out waiting for channel write".into()),
}
+ } else {
+ tx.send(OutboundMsg::Packet(packet.clone()))
+ .map_err(|e| format!("session gone: {e}"))?;
+ if let Ok(raw) = packet.encode_line() {
+ eprintln!(
+ "[kdeconnect] → {} (queued): {}",
+ packet.packet_type,
+ String::from_utf8_lossy(&raw).trim()
+ );
+ }
+ Ok(1)
}
- if ok == 0 {
- return Err(if last_err.is_empty() {
- format!("all channels to {peer_id} failed to write")
- } else {
- last_err
- });
- }
- Ok(ok)
}
/// Ping a trusted / live peer over the open KDE Connect channel.
///
/// KDE Connect ping is **one-way**: Valent shows a **desktop notification**
/// (title = device name “Global”, body = message). The Valent window itself
- /// does not change.
+ /// does not change. Requires **paired** + live channel; unpaired peers drop
+ /// plugin packets.
pub fn ping_peer(&self, peer_id: &str) -> String {
+ let trusted = self
+ .pair
+ .lock()
+ .map(|g| g.trust().is_trusted(peer_id))
+ .unwrap_or(false);
+ if !trusted {
+ let msg = format!(
+ "Ping skipped: not paired with {peer_id} yet — Accept pair first (Valent rejects unpaired plugins)"
+ );
+ if let Ok(mut a) = self.activity.lock() {
+ *a = msg.clone();
+ }
+ return msg;
+ }
// Use millisecond ids like official clients.
let pid = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as i64)
.unwrap_or_else(|_| self.next_packet_id());
let pkt = KdePacket::ping_with_message(pid, "Hello from Global");
- match self.send_to_peer(peer_id, pkt) {
- Ok(n) => {
- let msg = format!(
- "Ping written on {n} channel(s) — desktop notification “Global” / “Hello from Global”"
+ // Wait for the session thread to actually write (honest status).
+ match self.send_to_peer(peer_id, pkt, true) {
+ Ok(_) => {
+ let msg = String::from(
+ "Ping on wire — check desktop notification tray for “Global”",
);
if let Ok(mut a) = self.activity.lock() {
*a = msg.clone();
@@ -338,6 +383,18 @@ impl Engine {
}
}
+ /// Non-blocking ping for the UI thread (avoids Android ANR on TLS I/O).
+ pub fn ping_peer_async(self: &Arc, peer_id: String) {
+ let eng = Arc::clone(self);
+ if let Ok(mut a) = eng.activity.lock() {
+ *a = format!("Sending ping to {peer_id}…");
+ }
+ thread::spawn(move || {
+ let msg = eng.ping_peer(&peer_id);
+ eprintln!("[kdeconnect] async ping result: {msg}");
+ });
+ }
+
/// Apply pair effects including **SendPair over TCP+TLS**.
fn apply_effects(&self, effects: Vec) -> String {
let mut msgs = Vec::new();
@@ -348,7 +405,7 @@ impl Engine {
let pkt = KdePacket::pair_request(self.next_packet_id(), pair);
// Prefer live session; fall back to opening a new connection.
let r = self
- .send_to_peer(&peer_id, pkt)
+ .send_to_peer(&peer_id, pkt, false)
.map(|_| ())
.or_else(|_| self.send_pair_wire(&peer_id, pair));
match r {
@@ -493,6 +550,7 @@ fn accept_loop(
live_sessions: Arc,
verification_key: Arc>>,
linked_peer_name: Arc>>,
+ linked_peer_id: Arc>>,
sessions: Arc>>>,
session_gen: Arc,
activity: Arc>,
@@ -509,6 +567,7 @@ fn accept_loop(
let live_sessions = live_sessions.clone();
let verification_key = verification_key.clone();
let linked_peer_name = linked_peer_name.clone();
+ let linked_peer_id = linked_peer_id.clone();
let sessions = sessions.clone();
let session_gen = session_gen.clone();
let activity = activity.clone();
@@ -526,6 +585,7 @@ fn accept_loop(
live_sessions,
verification_key,
linked_peer_name,
+ linked_peer_id,
sessions,
session_gen,
activity,
@@ -563,6 +623,7 @@ fn handle_primary_session(
live_sessions: Arc,
verification_key: Arc>>,
linked_peer_name: Arc>>,
+ linked_peer_id: Arc>>,
sessions: Arc>>>,
session_gen: Arc,
activity: Arc>,
@@ -571,8 +632,9 @@ fn handle_primary_session(
set_ads_paused(&discovery, true);
eprintln!("[kdeconnect] primary channel up with {addr} (live={n})");
- // Shared stream: session thread reads; UI writes via send_to_peer.
+ // Shared stream: session thread reads + drains outbound; UI only enqueues.
let stream = Arc::new(Mutex::new(tls));
+ let (outbound_tx, outbound_rx) = std::sync::mpsc::channel::();
// Verification key from peer cert.
if let Ok(g) = stream.lock() {
@@ -613,22 +675,26 @@ fn handle_primary_session(
if let Ok(mut g) = linked_peer_name.lock() {
*g = Some(body.device_name.clone());
}
+ if let Ok(mut g) = linked_peer_id.lock() {
+ *g = Some(body.device_id.clone());
+ }
let gen = session_gen.fetch_add(1, Ordering::SeqCst);
let is_trusted = pair
.lock()
.map(|g| g.trust().is_trusted(&body.device_id))
.unwrap_or(false);
if let Ok(mut map) = sessions.lock() {
- map.entry(body.device_id.clone())
- .or_default()
- .push(SessionEntry {
- gen,
- stream: stream.clone(),
- });
+ // Valent opens many short-lived sockets; keep only this newest session
+ // so Ping/plugin traffic goes to the channel it is still reading.
+ let entry = SessionEntry {
+ gen,
+ stream: stream.clone(),
+ outbound: outbound_tx.clone(),
+ };
+ map.insert(body.device_id.clone(), vec![entry]);
registered = Some((body.device_id.clone(), gen));
- let n_sess = map.get(&body.device_id).map(|v| v.len()).unwrap_or(0);
eprintln!(
- "[kdeconnect] registered session gen={gen} for {} (now {n_sess} live)",
+ "[kdeconnect] registered sole session gen={gen} for {}",
body.device_id
);
}
@@ -637,7 +703,7 @@ fn handle_primary_session(
format!("connected to “{}” (paired) — use Ping below", body.device_name)
} else {
format!(
- "linked with “{}” — match verification key, then Request pairing in Valent",
+ "linked with “{}” — match verification key, then Accept pair request (or request from Valent)",
body.device_name
)
};
@@ -646,6 +712,10 @@ fn handle_primary_session(
}
loop {
+ // Drain UI-enqueued packets before (and after) each poll so Ping never
+ // waits on the stream mutex.
+ flush_outbound(&outbound_rx, &stream);
+
let read_result = {
let mut g = match stream.lock() {
Ok(g) => g,
@@ -654,7 +724,8 @@ fn handle_primary_session(
break;
}
};
- let _ = set_session_read_timeout(&mut g, Some(Duration::from_millis(200)));
+ // Short poll so outbound pings flush promptly (≤50ms).
+ let _ = set_session_read_timeout(&mut g, Some(Duration::from_millis(50)));
read_packet_line(&mut *g)
};
@@ -664,6 +735,7 @@ fn handle_primary_session(
"[kdeconnect] ← {} from {addr} type={}",
session_peer_name, pkt.packet_type
);
+ flush_outbound(&outbound_rx, &stream);
// Once Valent is actively using this channel, push battery + ping.
// Battery updates Valent’s device page (visible). Ping is notify-only.
@@ -765,37 +837,58 @@ fn handle_primary_session(
"[kdeconnect] pair packet pair={} from {peer_name} ({peer_id})",
pb.pair
);
+ if let Ok(mut g) = linked_peer_id.lock() {
+ *g = Some(peer_id.clone());
+ }
+ if let Ok(mut g) = linked_peer_name.lock() {
+ *g = Some(peer_name.clone());
+ }
+ // Incoming pair:true from Valent (not a confirm of our own request):
+ // auto-accept on the live stream so pairing completes without a
+ // second TCP hop. Also leave a clear status for the UI.
if pb.pair {
let st = pair
.lock()
.map(|g| g.state_of(&peer_id))
.unwrap_or(PairState::Idle);
- if st != PairState::OutgoingRequest && st != PairState::Paired {
+ if st != PairState::OutgoingRequest {
if let Ok(mut g) = pair.lock() {
- let _ = g.handle(PairEvent::ReceivedPair {
- peer_id: peer_id.clone(),
- peer_name: peer_name.clone(),
- peer_type: peer_type.clone(),
- body: pb.clone(),
- });
+ if st != PairState::Paired {
+ let _ = g.handle(PairEvent::ReceivedPair {
+ peer_id: peer_id.clone(),
+ peer_name: peer_name.clone(),
+ peer_type: peer_type.clone(),
+ body: pb.clone(),
+ });
+ }
let (effects, changed) = g.handle(PairEvent::AcceptIncoming {
peer_id: peer_id.clone(),
});
if changed {
let _ = g.trust().save(&trust_path);
}
+ let mut wrote_ok = false;
for e in effects {
match e {
PairEffect::SendPair { pair: p, .. } => {
- let reply = KdePacket::pair_request(2, p);
+ let reply = KdePacket::pair_request(
+ std::time::SystemTime::now()
+ .duration_since(std::time::UNIX_EPOCH)
+ .map(|d| d.as_millis() as i64)
+ .unwrap_or(2),
+ p,
+ );
if let Ok(mut sg) = stream.lock() {
match write_packet(&mut *sg, &reply) {
- Ok(()) => eprintln!(
- "[kdeconnect] → pair={p} reply to {peer_name}"
- ),
+ Ok(()) => {
+ wrote_ok = true;
+ eprintln!(
+ "[kdeconnect] → pair={p} accept to {peer_name}"
+ );
+ }
Err(err) => eprintln!(
- "[kdeconnect] write pair reply: {err}"
+ "[kdeconnect] write pair accept: {err}"
),
}
}
@@ -808,11 +901,27 @@ fn handle_primary_session(
}
}
}
+ let msg = if wrote_ok {
+ format!(
+ "Paired with “{peer_name}” — accept sent on channel"
+ )
+ } else {
+ format!(
+ "Pair request from “{peer_name}” — tap Accept if still unpaired"
+ )
+ };
+ if let Ok(mut st) = status.lock() {
+ *st = msg.clone();
+ }
+ if let Ok(mut a) = activity.lock() {
+ *a = msg;
+ }
}
continue;
}
}
+ // Unpair or confirmation of our OutgoingRequest.
if let Ok(mut g) = pair.lock() {
let (effects, changed) = g.handle(PairEvent::ReceivedPair {
peer_id: peer_id.clone(),
@@ -859,6 +968,7 @@ fn handle_primary_session(
}
}
+ // Unregister first so new pings fail fast; then fail any queued PacketAcks.
if let Some((id, gen)) = registered {
if let Ok(mut map) = sessions.lock() {
if let Some(list) = map.get_mut(&id) {
@@ -874,6 +984,12 @@ fn handle_primary_session(
}
}
}
+ drop(outbound_tx);
+ while let Ok(msg) = outbound_rx.try_recv() {
+ if let OutboundMsg::PacketAck(_, ack) = msg {
+ let _ = ack.send(Err("session closed before write".into()));
+ }
+ }
let left = live_sessions
.fetch_sub(1, Ordering::SeqCst)
.saturating_sub(1);
@@ -893,6 +1009,9 @@ fn handle_primary_session(
if let Ok(mut g) = linked_peer_name.lock() {
*g = None;
}
+ if let Ok(mut g) = linked_peer_id.lock() {
+ *g = None;
+ }
}
}
@@ -908,6 +1027,47 @@ fn set_session_read_timeout(
sock.set_read_timeout(timeout).map_err(|e| e.to_string())
}
+/// Write all queued packets while holding the stream briefly.
+fn flush_outbound(
+ rx: &std::sync::mpsc::Receiver,
+ stream: &Arc>,
+) {
+ loop {
+ match rx.try_recv() {
+ Ok(msg) => {
+ let (pkt, ack) = match msg {
+ OutboundMsg::Packet(p) => (p, None),
+ OutboundMsg::PacketAck(p, a) => (p, Some(a)),
+ };
+ let result = match stream.lock() {
+ Ok(mut g) => write_packet(&mut *g, &pkt).map_err(|e| e),
+ Err(e) => Err(e.to_string()),
+ };
+ match &result {
+ Ok(()) => {
+ if let Ok(raw) = pkt.encode_line() {
+ eprintln!(
+ "[kdeconnect] → {} (flushed): {}",
+ pkt.packet_type,
+ String::from_utf8_lossy(&raw).trim()
+ );
+ }
+ }
+ Err(e) => eprintln!(
+ "[kdeconnect] outbound write {} failed: {e}",
+ pkt.packet_type
+ ),
+ }
+ if let Some(ack) = ack {
+ let _ = ack.send(result);
+ }
+ }
+ Err(std::sync::mpsc::TryRecvError::Empty) => break,
+ Err(std::sync::mpsc::TryRecvError::Disconnected) => break,
+ }
+ }
+}
+
/// Process-wide engine for the device UI process.
static ENGINE: Mutex