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 struct InputLock(HANDLE);
impl InputLock {
pub fn acquire(timeout_ms: u32) -> Result<Self, CtlError> {
actl_core::stop::check_current()?;
let lock = Self::acquire_at(&actl_core::state::SignalPaths::default(), timeout_ms)?;
actl_core::stop::check_current()?;
Ok(lock)
}
pub fn acquire_at(
paths: &actl_core::state::SignalPaths,
timeout_ms: u32,
) -> Result<Self, CtlError> {
let invocation = actl_core::stop::Invocation::capture(paths);
invocation.check(paths)?;
unsafe {
let handle = CreateMutexExW(None, w!("actl-input-lock"), 0, 0x0010_0001)
.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 {
let lock = Self(handle);
invocation.check(paths)?;
crate::bg_style::recover_stale(paths);
Ok(lock)
} 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);
}
if done > 0 && actl_core::wait_control::pause_pending() {
return Err(CtlError::with_evidence(
actl_core::ErrorCode::Aborted,
"workflow input interrupted for pause",
serde_json::json!({
"stage":"handoff","reason":"user_paused",
"action_started":true,"delivered":done
}),
));
}
if let Err(mut error) = send(seg) {
let mut evidence = error
.evidence
.take()
.unwrap_or_else(|| serde_json::json!({}));
if !evidence.is_object() {
evidence = serde_json::json!({"cause":evidence});
}
evidence["delivered_before_error"] = done.into();
evidence["failed_segment_size"] = seg.len().into();
error.evidence = Some(evidence);
return Err(error);
}
done += seg.len();
}
Ok(done)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn segment_boundary_observes_pause_and_keeps_delivered_count() {
let dir = std::env::temp_dir().join(format!(
"actl-pause-seg-{}",
actl_core::snapshot::new_snapshot_id()
));
std::fs::create_dir_all(&dir).unwrap();
let marker = dir.join("pause-requested");
std::fs::write(&marker, b"1").unwrap();
let _scope = actl_core::wait_control::Scope::enter(marker.clone());
let items: Vec<u8> = std::iter::repeat_n(b'x', 64).collect();
let error = send_in_segments(&items, 32, |_| Ok(()), || true).unwrap_err();
assert_eq!(error.code, actl_core::ErrorCode::Aborted);
let evidence = error.evidence.unwrap();
assert_eq!(evidence["delivered"], 32, "{evidence}");
assert_eq!(evidence["reason"], "user_paused");
assert_eq!(evidence["stage"], "handoff");
std::fs::remove_file(&marker).unwrap();
let done = send_in_segments(&items, 32, |_| Ok(()), || true).unwrap();
assert_eq!(done, 64);
std::fs::remove_dir_all(&dir).unwrap();
}
#[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 isolated_paths = paths.clone();
let other = std::thread::spawn(move || match InputLock::acquire_at(&isolated_paths, 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_at(&paths, 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_preserves_delivered_prefix() {
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);
assert_eq!(err.evidence.as_ref().unwrap()["delivered_before_error"], 2);
}
}