actl-uia 0.1.8

Windows UIA backend: the ONLY crate allowed to touch COM/unsafe
//! Client of the companion's human handoff gate. No unattended bypass switch.
use actl_core::{
    CtlError,
    handoff::*,
    state::{SignalPaths, unix_ms},
};
use std::{
    cell::Cell,
    time::{Duration, Instant},
};
use windows::{
    Win32::{
        Foundation::{LPARAM, WPARAM},
        UI::WindowsAndMessaging::*,
    },
    core::w,
};
thread_local! { static ACTIVE: Cell<bool> = const { Cell::new(false) }; }
thread_local! { static FOREGROUND_READY: Cell<bool> = const { Cell::new(false) }; }
thread_local! { static FOREGROUND_DEPTH: Cell<usize> = const { Cell::new(0) }; }

/// Return eligibility to the companion before this short-lived caller exits.
pub(crate) struct ForegroundLease(u32);
fn release_outermost(depth: &Cell<usize>, release: impl FnOnce()) {
    let previous = depth.get();
    depth.set(previous.saturating_sub(1));
    if previous == 1 {
        release();
    }
}
impl Drop for ForegroundLease {
    fn drop(&mut self) {
        FOREGROUND_DEPTH.with(|depth| {
            release_outermost(depth, || {
                let result = unsafe { AllowSetForegroundWindow(self.0) };
                actl_core::signal_session::flow_event(
                    "foreground_return",
                    serde_json::json!({
                        "companion_pid":self.0,"accepted":result.is_ok(),
                        "failure":result.err().map(|e|e.to_string())
                    }),
                );
            })
        });
    }
}

pub(crate) fn foreground_lease() -> Result<Option<ForegroundLease>, CtlError> {
    if !ACTIVE.with(Cell::get) {
        return Ok(None);
    }
    check()?;
    if FOREGROUND_DEPTH.with(Cell::get) == 0
        && !FOREGROUND_READY.with(|ready| ready.replace(false))
        && query(FOREGROUND_MESSAGE)? != ALLOWED
    {
        return Err(failure("foreground_handoff_revoked"));
    }
    let hwnd = unsafe { FindWindowW(w!("actl_signal_window"), None) }
        .map_err(|_| failure("companion_missing"))?;
    let mut pid = 0;
    unsafe {
        GetWindowThreadProcessId(hwnd, Some(&mut pid));
    }
    if pid == 0
        || SignalPaths::default()
            .display_lease()
            .is_none_or(|lease| lease.pid != pid || !lease.compatible(unix_ms()))
    {
        return Err(failure("companion_identity_changed"));
    }
    FOREGROUND_DEPTH.with(|depth| depth.set(depth.get() + 1));
    Ok(Some(ForegroundLease(pid)))
}

fn failure(reason: &str) -> CtlError {
    actl_core::signal_session::handoff_event(reason);
    actl_core::handoff::error(reason)
}

fn query(message: u32) -> Result<usize, CtlError> {
    let hwnd = unsafe { FindWindowW(w!("actl_signal_window"), None) }
        .map_err(|_| failure("companion_missing"))?;
    let mut result = 0;
    let send = |result: &mut usize| unsafe {
        SendMessageTimeoutW(
            hwnd,
            message,
            WPARAM(std::process::id() as usize),
            LPARAM(0),
            SMTO_ABORTIFHUNG | SMTO_BLOCK,
            500,
            Some(result),
        )
    };
    let mut sent = send(&mut result);
    // Re-query consent once after a slow renderer response, never replay an action.
    if sent.0 == 0 && message == CHECK_MESSAGE {
        actl_core::stop::check_current()?;
        actl_core::signal_session::flow_event(
            "handoff_query_retry",
            serde_json::json!({"message":message,"timeout_ms":500}),
        );
        result = 0;
        sent = send(&mut result);
    }
    if sent.0 == 0 || result == 0 {
        let mut error = failure("companion_unresponsive");
        if let Some(evidence) = error.evidence.as_mut() {
            evidence["query_message"] = message.into();
            evidence["query_timeout_ms"] = 500.into();
            evidence["query_return"] = result.into();
        }
        return Err(error);
    }
    Ok(result)
}
pub(crate) fn check() -> Result<(), CtlError> {
    if ACTIVE.with(Cell::get) {
        let started = std::time::Instant::now();
        let mut code = query(CHECK_MESSAGE)?;
        if code == actl_core::handoff::YIELD_BUSY {
            // 批次3(显示协议 v11):Standing 授权仍在、输入活跃——有界等待安静
            // (≤3s,150ms 轮询),期间急停/停止即时生效;超时才按让行失败。
            // 误触(鼠标一晃/键盘一碰)不再直接打断任务。
            let deadline = started + std::time::Duration::from_millis(3000);
            loop {
                std::thread::sleep(std::time::Duration::from_millis(150));
                actl_core::stop::check_current()?;
                code = query(CHECK_MESSAGE)?;
                if code != actl_core::handoff::YIELD_BUSY || std::time::Instant::now() >= deadline {
                    break;
                }
            }
        }
        match code {
            ALLOWED => {}
            USER_PAUSED => return Err(failure("user_paused")),
            CANCELLED => {
                return Err(CtlError::new(
                    actl_core::ErrorCode::Aborted,
                    "execution generation changed",
                ));
            }
            MONITOR_UNAVAILABLE => return Err(failure("input_monitor_unavailable")),
            _ => {
                let mut error = failure("human_takeover_or_grant_expired");
                let waited = started.elapsed().as_millis() as u64;
                if waited > 0
                    && let Some(evidence) = error.evidence.as_mut()
                {
                    evidence["yield"] = "busy_timeout".into();
                    evidence["waited_ms"] = waited.into();
                }
                return Err(error);
            }
        }
    }
    Ok(())
}
pub(crate) fn prepare(target: &str, action: &str) -> Result<(), CtlError> {
    if ACTIVE.with(Cell::get) {
        return check();
    }
    request(target, action).map_err(|mut error| {
        if let Some(evidence) = error.evidence.as_mut()
            && evidence["stage"] == "handoff"
        {
            evidence["action_started"] = false.into();
        }
        error
    })
}
fn request(target: &str, action: &str) -> Result<(), CtlError> {
    let (call, task) =
        actl_core::signal_session::handoff_identity().ok_or_else(|| failure("missing_session"))?;
    let paths = SignalPaths::default();
    let file = format!("handoff-{}.json", std::process::id());
    struct Remove(std::path::PathBuf);
    impl Drop for Remove {
        fn drop(&mut self) {
            let _ = std::fs::remove_file(&self.0);
        }
    }
    let _remove = Remove(paths.dir.join(&file));
    let mut request = Request {
        approval: actl_core::handoff::configured(),
        generation: actl_core::stop::generation()?,
        call,
        task,
        target: target.chars().take(180).collect(),
        ts_ms: unix_ms(),
        action: action.chars().take(32).collect(),
    };
    let mut waiting = false;
    let start = Instant::now();
    while start.elapsed() < Duration::from_secs(25) {
        crate::feedback::check_stop()?;
        request.ts_ms = unix_ms();
        paths
            .write_display(&file, &request)
            .map_err(crate::internal)?;
        match query(REQUEST_MESSAGE)? {
            ALLOWED => {
                actl_core::signal_session::handoff_event("allowed");
                ACTIVE.with(|a| a.set(true));
                FOREGROUND_READY.with(|ready| ready.set(true));
                return Ok(());
            }
            PANEL_UNAVAILABLE => return Err(failure("handoff_panel_unavailable")),
            MONITOR_UNAVAILABLE => return Err(failure("input_monitor_unavailable")),
            USER_PAUSED => return Err(failure("user_paused")),
            DEFERRED => return Err(failure("user_deferred")),
            CANCELLED => {
                return Err(CtlError::new(
                    actl_core::ErrorCode::Aborted,
                    "execution generation changed",
                ));
            }
            PENDING => {
                if !waiting {
                    actl_core::signal_session::handoff_event("handoff_pending");
                    waiting = true;
                }
                std::thread::sleep(Duration::from_millis(crate::timing().handoff_poll_ms));
            }
            _ => return Err(failure("request_rejected")),
        }
    }
    Err(failure("handoff_pending"))
}

#[cfg(test)]
mod foreground_tests {
    use super::*;
    #[test]
    fn nested_actions_return_eligibility_only_at_outer_boundary() {
        let depth = Cell::new(2);
        let returns = Cell::new(0);
        release_outermost(&depth, || returns.set(returns.get() + 1));
        assert_eq!(returns.get(), 0);
        release_outermost(&depth, || returns.set(returns.get() + 1));
        assert_eq!(returns.get(), 1);
        release_outermost(&depth, || panic!("no lease to return"));
        depth.set(1); // Next action in the same process gets its own lease.
        release_outermost(&depth, || returns.set(returns.get() + 1));
        assert_eq!(returns.get(), 2);
    }
}