use crate::clipboard::Clipboard as _;
use crate::config::PasteShortcut;
use crate::error::TalkError;
use crate::paste::node::{PasteCtx, PasteNode};
use crate::paste::{log_preview, simulate_paste};
use async_trait::async_trait;
use std::sync::atomic::Ordering;
use std::time::{Duration, Instant};
const GATE_POLL_INTERVAL_MS: u64 = 5;
#[derive(Debug, Clone, Copy)]
pub(crate) struct ClipboardNode {
pub(crate) shortcut: PasteShortcut,
#[allow(dead_code)] pub(crate) restore_settle_ms: u64,
pub(crate) chunk_fetch_timeout_ms: u64,
pub(crate) target_quiescence_ms: u64,
pub(crate) target_fetch_retries: u32,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum GateDecision {
Learned { expected: u32 },
Confirmed { observed: u32 },
AbortedTimeout { observed: u32, required: u32 },
}
#[async_trait]
impl PasteNode for ClipboardNode {
async fn paste(&self, text: &str, ctx: &PasteCtx<'_>) -> Result<(), TalkError> {
let clipboard = ctx.clipboard;
let target_base = match ctx.target_client_base {
Some(base) => base,
None => {
self.serve_and_simulate(clipboard, text).await?;
self.run_fallback_gate(clipboard).await;
let _ = ctx.t_stop;
return Ok(());
}
};
let learn_phase = ctx.expected_target_fetches.load(Ordering::Relaxed) == 0;
for attempt in 0..=self.target_fetch_retries {
if attempt > 0 {
if learn_phase {
ctx.expected_target_fetches.store(0, Ordering::Relaxed);
}
if let Some(wid) = ctx.target_window {
if let Err(e) = crate::paste::ensure_focus(wid).await {
log::warn!(
"paste(clipboard-node): retry {} could not re-focus \
target window {}: {} — retrying anyway",
attempt,
wid,
e,
);
}
}
}
self.serve_and_simulate(clipboard, text).await?;
match self.run_target_gate(clipboard, target_base, ctx).await {
Ok(()) => {
if attempt > 0 {
log::info!(
"paste(clipboard-node): chunk confirmed on retry {} \
(target client-base {:#x})",
attempt,
target_base,
);
}
let _ = ctx.t_stop;
return Ok(());
}
Err(e) => {
if attempt < self.target_fetch_retries {
log::warn!(
"paste(clipboard-node): target client-base {:#x} did not \
fetch chunk within {} ms (attempt {}/{}) — re-focusing \
and retrying: {}",
target_base,
self.chunk_fetch_timeout_ms,
attempt + 1,
self.target_fetch_retries + 1,
e,
);
continue;
}
self.signal_final_abort(ctx, &e);
return Err(e);
}
}
}
Err(TalkError::Clipboard(
"paste aborted: retry loop exited without a decision".to_string(),
))
}
}
impl ClipboardNode {
async fn serve_and_simulate(
&self,
clipboard: &crate::clipboard::X11Clipboard,
text: &str,
) -> Result<(), TalkError> {
log::trace!(
"paste(clipboard-node): set chunk content={}",
log_preview(text),
);
clipboard.set_text(text).await?;
if log::log_enabled!(log::Level::Trace) {
match clipboard.get_text().await {
Ok(rb) if rb == text => {
log::trace!("paste(clipboard-node): chunk read-back OK");
}
Ok(rb) => {
log::trace!(
"paste(clipboard-node): read-back MISMATCH — clipboard holds {}",
log_preview(&rb),
);
}
Err(e) => {
log::trace!("paste(clipboard-node): read-back failed: {}", e);
}
}
}
tokio::time::sleep(Duration::from_millis(5)).await;
simulate_paste(self.shortcut).await
}
fn signal_final_abort(&self, ctx: &PasteCtx<'_>, err: &TalkError) {
ctx.sink.emit(crate::telemetry::TranscriptionEvent::Failed {
reason: err.to_string(),
t: Instant::now(),
});
if let Some(alert) = ctx.alert.as_ref() {
alert();
}
}
async fn run_target_gate(
&self,
clipboard: &crate::clipboard::X11Clipboard,
target_base: u32,
ctx: &PasteCtx<'_>,
) -> Result<(), TalkError> {
let expected_prev = ctx.expected_target_fetches.load(Ordering::Relaxed);
let timeout = Duration::from_millis(self.chunk_fetch_timeout_ms);
let quiescence = Duration::from_millis(self.target_quiescence_ms);
let decision = if expected_prev == 0 {
wait_and_learn(clipboard, target_base, timeout, quiescence).await
} else {
wait_and_confirm(clipboard, target_base, expected_prev, timeout, quiescence).await
};
match decision {
GateDecision::Learned { expected } => {
ctx.expected_target_fetches
.store(expected, Ordering::Relaxed);
log::info!(
"paste(clipboard-node): target client-base {:#x} learned \
expected_target_fetches={} (chunk 1 quiesced after {} ms)",
target_base,
expected,
self.target_quiescence_ms,
);
Ok(())
}
GateDecision::Confirmed { observed } => {
log::trace!(
"paste(clipboard-node): target client-base {:#x} confirmed \
fetches={} (>= expected={})",
target_base,
observed,
expected_prev,
);
Ok(())
}
GateDecision::AbortedTimeout { observed, required } => {
let msg = if expected_prev == 0 {
format!(
"paste aborted: target X11 client-base {:#x} never fetched \
the clipboard for chunk 1 within {} ms (observed={}) — \
wrong focus, unsupported app, or shortcut mismatch",
target_base, self.chunk_fetch_timeout_ms, observed,
)
} else {
format!(
"paste aborted: target X11 client-base {:#x} only fetched \
clipboard {}/{} times within {} ms (this chunk would be \
dropped — refusing to overwrite silently)",
target_base, observed, required, self.chunk_fetch_timeout_ms,
)
};
log::error!("{}", msg);
Err(TalkError::Clipboard(msg))
}
}
}
async fn run_fallback_gate(&self, clipboard: &crate::clipboard::X11Clipboard) {
log::debug!(
"paste(clipboard-node): no target client-base — falling back to \
legacy served_count gate (chunk_fetch_timeout_ms={})",
self.chunk_fetch_timeout_ms,
);
let served = clipboard
.wait_until_served(0, Duration::from_millis(self.chunk_fetch_timeout_ms))
.await;
if served == 0 {
log::warn!(
"paste(clipboard-node): blind-paste fallback timed out after {} ms \
with served_count=0 — target never fetched our clipboard, \
likely paste corruption",
self.chunk_fetch_timeout_ms,
);
} else {
log::trace!(
"paste(clipboard-node): blind-paste consumed (served_count={})",
served,
);
}
}
}
async fn wait_and_learn(
clipboard: &crate::clipboard::X11Clipboard,
target_base: u32,
timeout: Duration,
quiescence: Duration,
) -> GateDecision {
let deadline = Instant::now() + timeout;
let poll = Duration::from_millis(GATE_POLL_INTERVAL_MS);
loop {
let count = clipboard.target_fetch_count(target_base);
if count > 0 {
break;
}
if Instant::now() >= deadline {
return GateDecision::AbortedTimeout {
observed: 0,
required: 1,
};
}
tokio::time::sleep(poll).await;
}
let mut last_count = clipboard.target_fetch_count(target_base);
let mut last_change = Instant::now();
loop {
if last_change.elapsed() >= quiescence {
return GateDecision::Learned {
expected: last_count,
};
}
if Instant::now() >= deadline {
return GateDecision::Learned {
expected: last_count,
};
}
tokio::time::sleep(poll).await;
let now = clipboard.target_fetch_count(target_base);
if now != last_count {
last_count = now;
last_change = Instant::now();
}
}
}
async fn wait_and_confirm(
clipboard: &crate::clipboard::X11Clipboard,
target_base: u32,
expected: u32,
timeout: Duration,
quiescence: Duration,
) -> GateDecision {
let deadline = Instant::now() + timeout;
let poll = Duration::from_millis(GATE_POLL_INTERVAL_MS);
let mut observed;
loop {
observed = clipboard.target_fetch_count(target_base);
if observed >= expected {
break;
}
if Instant::now() >= deadline {
return GateDecision::AbortedTimeout {
observed,
required: expected,
};
}
tokio::time::sleep(poll).await;
}
let mut last_count = observed;
let mut last_change = Instant::now();
loop {
if last_change.elapsed() >= quiescence {
return GateDecision::Confirmed {
observed: last_count,
};
}
if Instant::now() >= deadline {
return GateDecision::Confirmed {
observed: last_count,
};
}
tokio::time::sleep(poll).await;
let now = clipboard.target_fetch_count(target_base);
if now != last_count {
last_count = now;
last_change = Instant::now();
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::x11::clipboard::client_base;
#[test]
fn client_base_groups_target_and_child_widget() {
let mask: u32 = 0x001F_FFFF;
let target: u32 = 50_331_661;
let child: u32 = 50_331_792;
let expected_base: u32 = 0x0300_0000;
assert_eq!(client_base(target, mask), expected_base);
assert_eq!(client_base(child, mask), expected_base);
}
#[test]
fn client_base_distinguishes_clipboard_managers_from_target() {
let mask: u32 = 0x001F_FFFF;
let target_base = client_base(50_331_661, mask);
let manager_a: u32 = 0x0640_0001;
let manager_b: u32 = 0x0680_1234;
assert_ne!(client_base(manager_a, mask), target_base);
assert_ne!(client_base(manager_b, mask), target_base);
}
#[test]
fn client_base_is_pure_bit_masking() {
let mask: u32 = 0x0000_FFFF;
assert_eq!(client_base(0x1234_5678, mask), 0x1234_0000);
assert_eq!(client_base(0x1234_FFFF, mask), 0x1234_0000);
let mask: u32 = 0x0000_00FF;
assert_eq!(client_base(0xABCD_EF12, mask), 0xABCD_EF00);
}
#[tokio::test]
async fn learn_aborts_when_target_never_fetches() {
let timeout = Duration::from_millis(20);
let quiescence = Duration::from_millis(50);
let decision = simulate_learn_with(|_| 0, timeout, quiescence).await;
match decision {
GateDecision::AbortedTimeout { observed, required } => {
assert_eq!(observed, 0);
assert_eq!(required, 1);
}
other => panic!("expected AbortedTimeout, got {:?}", other),
}
}
#[tokio::test]
async fn learn_freezes_at_one_when_only_one_fetch_arrives() {
let started = Instant::now();
let decision = simulate_learn_with(
move |_| {
if started.elapsed() > Duration::from_millis(2) {
1
} else {
0
}
},
Duration::from_millis(300),
Duration::from_millis(40),
)
.await;
match decision {
GateDecision::Learned { expected } => assert_eq!(expected, 1),
other => panic!("expected Learned{{1}}, got {:?}", other),
}
}
#[tokio::test]
async fn learn_freezes_at_two_when_target_fetches_twice() {
let started = Instant::now();
let decision = simulate_learn_with(
move |_| {
let e = started.elapsed();
if e > Duration::from_millis(15) {
2
} else if e > Duration::from_millis(2) {
1
} else {
0
}
},
Duration::from_millis(300),
Duration::from_millis(40),
)
.await;
match decision {
GateDecision::Learned { expected } => assert_eq!(expected, 2),
other => panic!("expected Learned{{2}}, got {:?}", other),
}
}
#[tokio::test]
async fn confirm_succeeds_when_expected_count_reached() {
let started = Instant::now();
let decision = simulate_confirm_with(
move |_| {
let e = started.elapsed();
if e > Duration::from_millis(15) {
2
} else if e > Duration::from_millis(2) {
1
} else {
0
}
},
2,
Duration::from_millis(300),
Duration::from_millis(40),
)
.await;
match decision {
GateDecision::Confirmed { observed } => assert!(observed >= 2),
other => panic!("expected Confirmed, got {:?}", other),
}
}
#[tokio::test]
async fn confirm_aborts_when_expected_count_not_reached() {
let started = Instant::now();
let decision = simulate_confirm_with(
move |_| {
if started.elapsed() > Duration::from_millis(2) {
1
} else {
0
}
},
2,
Duration::from_millis(40),
Duration::from_millis(20),
)
.await;
match decision {
GateDecision::AbortedTimeout { observed, required } => {
assert_eq!(observed, 1);
assert_eq!(required, 2);
}
other => panic!("expected AbortedTimeout, got {:?}", other),
}
}
type AttemptResult = Result<(), String>;
struct RetryTrace {
result: AttemptResult,
resets: u32,
signalled: bool,
attempts: u32,
}
fn run_retry_plan<F>(retries: u32, learn_phase: bool, mut gate: F) -> RetryTrace
where
F: FnMut(u32) -> AttemptResult,
{
let mut resets = 0u32;
let mut attempts = 0u32;
for attempt in 0..=retries {
if attempt > 0 && learn_phase {
resets += 1;
}
attempts += 1;
match gate(attempt) {
Ok(()) => {
return RetryTrace {
result: Ok(()),
resets,
signalled: false,
attempts,
};
}
Err(e) => {
if attempt < retries {
continue;
}
return RetryTrace {
result: Err(e),
resets,
signalled: true,
attempts,
};
}
}
}
RetryTrace {
result: Err("retry loop exited without a decision".to_string()),
resets,
signalled: false,
attempts,
}
}
#[test]
fn retry_succeeds_after_one_failed_attempt() {
let trace = run_retry_plan(2, true, |attempt| {
if attempt == 0 {
Err("first attempt: target never fetched".to_string())
} else {
Ok(())
}
});
assert!(trace.result.is_ok(), "second attempt should confirm");
assert!(!trace.signalled, "no abort signal on eventual success");
assert_eq!(trace.attempts, 2, "one failure + one success");
assert_eq!(trace.resets, 1);
}
#[test]
fn retry_aborts_and_signals_once_after_exhaustion() {
let mut fail_count = 0u32;
let trace = run_retry_plan(2, false, |_attempt| {
fail_count += 1;
Err("target never fetched".to_string())
});
assert!(trace.result.is_err(), "exhausted retries must abort");
assert!(
trace.signalled,
"abort signal fires exactly once on final abort"
);
assert_eq!(trace.attempts, 3, "1 initial + 2 retries = 3 attempts");
assert_eq!(fail_count, 3, "gate invoked once per attempt");
assert_eq!(trace.resets, 0);
}
#[test]
fn chunk_one_retries_relearn_but_chunk_n_does_not() {
let learn = run_retry_plan(2, true, |_| Err("nope".to_string()));
assert_eq!(learn.resets, 2, "chunk-1 relearns on each of 2 retries");
let confirm = run_retry_plan(2, false, |_| Err("nope".to_string()));
assert_eq!(confirm.resets, 0, "chunk-N never resets learned count");
}
#[test]
fn zero_retries_aborts_on_first_failure() {
let trace = run_retry_plan(0, true, |_| Err("nope".to_string()));
assert!(trace.result.is_err());
assert!(trace.signalled);
assert_eq!(trace.attempts, 1);
}
#[test]
fn first_attempt_success_no_retry_no_signal() {
let trace = run_retry_plan(2, true, |attempt| {
assert_eq!(attempt, 0, "must not run a second attempt");
Ok(())
});
assert!(trace.result.is_ok());
assert!(!trace.signalled);
assert_eq!(trace.attempts, 1);
assert_eq!(trace.resets, 0);
}
async fn simulate_learn_with<F>(
count: F,
timeout: Duration,
quiescence: Duration,
) -> GateDecision
where
F: Fn(Instant) -> u32,
{
let deadline = Instant::now() + timeout;
let poll = Duration::from_millis(GATE_POLL_INTERVAL_MS);
loop {
let c = count(Instant::now());
if c > 0 {
break;
}
if Instant::now() >= deadline {
return GateDecision::AbortedTimeout {
observed: 0,
required: 1,
};
}
tokio::time::sleep(poll).await;
}
let mut last_count = count(Instant::now());
let mut last_change = Instant::now();
loop {
if last_change.elapsed() >= quiescence {
return GateDecision::Learned {
expected: last_count,
};
}
if Instant::now() >= deadline {
return GateDecision::Learned {
expected: last_count,
};
}
tokio::time::sleep(poll).await;
let now = count(Instant::now());
if now != last_count {
last_count = now;
last_change = Instant::now();
}
}
}
async fn simulate_confirm_with<F>(
count: F,
expected: u32,
timeout: Duration,
quiescence: Duration,
) -> GateDecision
where
F: Fn(Instant) -> u32,
{
let deadline = Instant::now() + timeout;
let poll = Duration::from_millis(GATE_POLL_INTERVAL_MS);
let mut observed;
loop {
observed = count(Instant::now());
if observed >= expected {
break;
}
if Instant::now() >= deadline {
return GateDecision::AbortedTimeout {
observed,
required: expected,
};
}
tokio::time::sleep(poll).await;
}
let mut last_count = observed;
let mut last_change = Instant::now();
loop {
if last_change.elapsed() >= quiescence {
return GateDecision::Confirmed {
observed: last_count,
};
}
if Instant::now() >= deadline {
return GateDecision::Confirmed {
observed: last_count,
};
}
tokio::time::sleep(poll).await;
let now = count(Instant::now());
if now != last_count {
last_count = now;
last_change = Instant::now();
}
}
}
}