use actl_core::{CtlError, ErrorCode};
use windows::Win32::Foundation::{
CloseHandle, HANDLE, WAIT_ABANDONED, WAIT_OBJECT_0, WAIT_TIMEOUT,
};
use windows::Win32::System::Threading::{CreateMutexExW, ReleaseMutex, WaitForSingleObject};
use windows::core::w;
pub const LOCK_TIMEOUT_MS: u32 = 3000;
pub struct InputLock(HANDLE);
impl InputLock {
pub fn acquire(timeout_ms: u32) -> Result<Self, CtlError> {
Self::acquire_at(&actl_core::state::SignalPaths::default(), timeout_ms)
}
pub fn acquire_at(
paths: &actl_core::state::SignalPaths,
timeout_ms: u32,
) -> Result<Self, CtlError> {
if paths.stop_requested() {
return Err(CtlError::new(
ErrorCode::Aborted,
"stop requested via signal channel (emergency hotkey); clear the flag to resume",
));
}
unsafe {
let handle = CreateMutexExW(None, w!("actl-input-lock"), 0, 0x0010_0000)
.map_err(|e| CtlError::internal(format!("CreateMutexExW(actl-input-lock): {e}")))?;
let waited = WaitForSingleObject(handle, timeout_ms);
if waited == WAIT_OBJECT_0 || waited == WAIT_ABANDONED {
Ok(Self(handle))
} else if waited == WAIT_TIMEOUT {
let _ = CloseHandle(handle);
Err(CtlError::new(
ErrorCode::Timeout,
format!(
"another actl instance holds the input lock (waited {timeout_ms}ms); \
retry after it finishes or raise --timeout"
),
))
} else {
let _ = CloseHandle(handle);
Err(CtlError::internal(format!(
"WaitForSingleObject returned {waited:?}"
)))
}
}
}
}
impl Drop for InputLock {
fn drop(&mut self) {
unsafe {
let _ = ReleaseMutex(self.0);
let _ = CloseHandle(self.0);
}
}
}
pub fn send_in_segments<T, F, V>(
items: &[T],
segment_len: usize,
mut send: F,
mut focus_ok: V,
) -> Result<usize, CtlError>
where
F: FnMut(&[T]) -> Result<(), CtlError>,
V: FnMut() -> bool,
{
let segment_len = segment_len.max(1);
let mut done = 0usize;
for seg in items.chunks(segment_len) {
if done > 0 && !focus_ok() {
return Ok(done);
}
send(seg)?;
done += seg.len();
}
Ok(done)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn lock_is_mutual_across_threads() {
let paths = actl_core::state::SignalPaths::at(
std::env::temp_dir().join(format!("actl-lock-test-{}", std::process::id())),
);
let a = InputLock::acquire_at(&paths, 0).expect("first acquire");
let other = std::thread::spawn(|| match InputLock::acquire(0) {
Err(e) => Some(e.code),
Ok(_) => None,
});
match other.join().expect("thread ok") {
Some(code) => assert_eq!(code, ErrorCode::Timeout),
None => panic!("second acquire must time out"),
}
drop(a);
assert!(InputLock::acquire(0).is_ok());
}
#[test]
fn segments_send_all_when_focus_holds() {
let items: Vec<usize> = (0..70).collect();
let mut sent: Vec<usize> = Vec::new();
let done = send_in_segments(
&items,
32,
|seg| {
sent.extend_from_slice(seg);
Ok(())
},
|| true,
)
.unwrap();
assert_eq!(done, 70);
assert_eq!(sent, items);
}
#[test]
fn segments_stop_at_boundary_when_focus_lost() {
let items: Vec<usize> = (0..100).collect();
let mut calls = 0;
let done = send_in_segments(
&items,
10,
|_seg| Ok(()),
|| {
calls += 1;
calls < 2
},
)
.unwrap();
assert_eq!(done, 20, "停在批边界,报告已送达量");
}
#[test]
fn segment_error_propagates_with_done_loss() {
let items: Vec<usize> = vec![1, 2, 3, 4];
let err = send_in_segments(
&items,
2,
|seg| {
if seg[0] == 3 {
Err(CtlError::internal("inject failed"))
} else {
Ok(())
}
},
|| true,
)
.unwrap_err();
assert_eq!(err.code, ErrorCode::Internal);
}
}