mod capture;
pub mod client;
pub mod protocol;
pub mod server;
pub mod session;
use chrono::TimeZone;
use protocol::{
EvaluateResult, HistoryEntry, HistoryParams, HistoryResult, INPUT_ALREADY_PENDING, IpcMethod,
IpcRequest, IpcResponse, R_BUSY, R_NOT_AT_PROMPT, RSessionInfo, SessionResult, USER_IS_TYPING,
UserInputResult,
};
use std::path::PathBuf;
use std::sync::{
Arc, Mutex, OnceLock,
atomic::{AtomicBool, Ordering},
};
static HEADLESS_SHUTDOWN: OnceLock<Arc<AtomicBool>> = OnceLock::new();
static HEADLESS_HISTORY: OnceLock<Mutex<reedline::SqliteBackedHistory>> = OnceLock::new();
static HISTORY_DB_INFO: OnceLock<(PathBuf, Option<reedline::HistorySessionId>)> = OnceLock::new();
static IPC_RECEIVER: OnceLock<Mutex<Option<std::sync::mpsc::Receiver<IpcRequest>>>> =
OnceLock::new();
static PENDING_IPC_OPERATION: OnceLock<Mutex<Option<PendingIpcOperation>>> = OnceLock::new();
fn pending_ipc_operation() -> &'static Mutex<Option<PendingIpcOperation>> {
PENDING_IPC_OPERATION.get_or_init(|| Mutex::new(None))
}
static R_IS_AT_PROMPT: OnceLock<AtomicBool> = OnceLock::new();
static IN_ALTERNATE_MODE: AtomicBool = AtomicBool::new(false);
fn r_is_at_prompt() -> &'static AtomicBool {
R_IS_AT_PROMPT.get_or_init(|| AtomicBool::new(false))
}
pub fn set_in_alternate_mode(active: bool) {
IN_ALTERNATE_MODE.store(active, Ordering::Release);
}
pub fn is_in_alternate_mode() -> bool {
IN_ALTERNATE_MODE.load(Ordering::Acquire)
}
static BREAK_SIGNAL: OnceLock<Arc<AtomicBool>> = OnceLock::new();
static PENDING_VISIBLE_EVAL: OnceLock<Mutex<Option<PendingVisibleEval>>> = OnceLock::new();
struct PendingVisibleEval {
reply: tokio::sync::oneshot::Sender<IpcResponse>,
deadline: std::time::Instant,
timeout: std::time::Duration,
}
fn pending_visible_eval() -> &'static Mutex<Option<PendingVisibleEval>> {
PENDING_VISIBLE_EVAL.get_or_init(|| Mutex::new(None))
}
pub struct PendingIpcOperation {
pub kind: PendingIpcKind,
pub code: String,
}
pub enum PendingIpcKind {
SilentEvaluate {
reply: tokio::sync::oneshot::Sender<IpcResponse>,
},
VisibleEvaluate {
reply: tokio::sync::oneshot::Sender<IpcResponse>,
timeout: std::time::Duration,
},
UserInput {
reply: tokio::sync::oneshot::Sender<IpcResponse>,
},
}
pub fn break_signal() -> Arc<AtomicBool> {
BREAK_SIGNAL
.get_or_init(|| Arc::new(AtomicBool::new(false)))
.clone()
}
pub fn start_server(
bind: Option<&str>,
log_file: Option<String>,
history_session_id: Option<i64>,
) -> std::io::Result<session::SessionInfo> {
let (tx, rx) = std::sync::mpsc::channel();
let _ = pending_ipc_operation();
let started_at = chrono::Local::now().to_rfc3339();
let session = server::start_server(tx, bind, &started_at, log_file, history_session_id)?;
match IPC_RECEIVER.get() {
Some(existing) => {
*existing.lock().unwrap() = Some(rx);
}
None => {
IPC_RECEIVER
.set(Mutex::new(Some(rx)))
.map_err(|_| std::io::Error::other("IPC receiver already initialized"))?;
}
}
Ok(session)
}
pub fn stop_server() {
if let Some(receiver) = IPC_RECEIVER.get() {
*receiver.lock().unwrap() = None;
}
if let Some(pending) = take_pending_ipc_operation() {
match pending.kind {
PendingIpcKind::SilentEvaluate { reply }
| PendingIpcKind::VisibleEvaluate { reply, .. }
| PendingIpcKind::UserInput { reply } => {
let _ = reply.send(IpcResponse::error(
R_NOT_AT_PROMPT,
"IPC server is shutting down".to_string(),
));
}
}
}
if let Some(pending) = pending_visible_eval()
.lock()
.unwrap_or_else(|e| e.into_inner())
.take()
{
let (stdout, stderr) = arf_libr::finish_ipc_capture();
let _ = pending.reply.send(IpcResponse::error(
R_NOT_AT_PROMPT,
format!(
"IPC server shut down during visible evaluate (stdout: {} bytes, stderr: {} bytes)",
stdout.len(),
stderr.len(),
),
));
}
server::stop_server();
}
pub fn poll_ipc_requests() {
check_visible_eval_completion();
let receiver = match IPC_RECEIVER.get() {
Some(r) => r,
None => return, };
let rx = match receiver.try_lock() {
Ok(rx) => rx,
Err(_) => return,
};
let rx = match rx.as_ref() {
Some(rx) => rx,
None => return, };
while let Ok(request) = rx.try_recv() {
handle_request(request);
}
}
pub(in crate::ipc) const DEFAULT_EVAL_TIMEOUT: std::time::Duration =
std::time::Duration::from_secs(300);
fn check_visible_eval_completion() {
let at_prompt = r_is_at_prompt().load(Ordering::Acquire);
let mut guard = match pending_visible_eval().try_lock() {
Ok(g) => g,
Err(_) => return,
};
if let Some(pending) = guard.as_ref()
&& std::time::Instant::now() > pending.deadline
{
let pending = guard.take().unwrap();
let (stdout, stderr) = arf_libr::finish_ipc_capture();
let _ = pending.reply.send(IpcResponse::error(
R_BUSY,
format!(
"Visible evaluate timed out after {}s (stdout: {} bytes, stderr: {} bytes)",
pending.timeout.as_secs(),
stdout.len(),
stderr.len(),
),
));
return;
}
if !at_prompt {
return;
}
if let Some(pending) = guard.take() {
r_is_at_prompt().store(false, Ordering::Release);
let (stdout, stderr) = arf_libr::finish_ipc_capture();
let result = EvaluateResult {
stdout,
stderr,
value: None,
error: None,
};
let _ = pending.reply.send(IpcResponse::Evaluate(result));
r_is_at_prompt().store(true, Ordering::Release);
}
}
fn handle_request(request: IpcRequest) {
let IpcRequest { method, reply } = request;
if matches!(method, IpcMethod::Session) {
let r_at_prompt = r_is_at_prompt().load(Ordering::Acquire);
let has_pending = pending_ipc_operation()
.lock()
.unwrap_or_else(|e| e.into_inner())
.is_some();
let in_alt_mode = is_in_alternate_mode();
let try_r = r_at_prompt && !has_pending && !in_alt_mode;
let reason = if in_alt_mode {
"R is in alternate mode (shell, history browser, or help browser)"
} else if !r_at_prompt {
"R is busy evaluating another expression"
} else if has_pending {
"Another IPC operation is pending"
} else {
"" };
let result = collect_session_result(try_r, reason);
let _ = reply.send(IpcResponse::Session(Box::new(result)));
return;
}
if is_in_alternate_mode() {
let _ = reply.send(IpcResponse::error(
R_NOT_AT_PROMPT,
"R is not at the command prompt".to_string(),
));
return;
}
if pending_ipc_operation()
.lock()
.unwrap_or_else(|e| e.into_inner())
.is_some()
{
let _ = reply.send(IpcResponse::error(
INPUT_ALREADY_PENDING,
"Another IPC operation is pending".to_string(),
));
return;
}
match method {
IpcMethod::Evaluate {
code,
visible,
timeout_ms,
} => {
if !r_is_at_prompt().load(Ordering::Acquire) {
let _ = reply.send(IpcResponse::error(R_BUSY, "R is busy".to_string()));
return;
}
let timeout = timeout_ms
.map(std::time::Duration::from_millis)
.unwrap_or(DEFAULT_EVAL_TIMEOUT);
let kind = if visible {
PendingIpcKind::VisibleEvaluate { reply, timeout }
} else {
PendingIpcKind::SilentEvaluate { reply }
};
*pending_ipc_operation()
.lock()
.unwrap_or_else(|e| e.into_inner()) = Some(PendingIpcOperation { kind, code });
fire_break_signal();
}
IpcMethod::UserInput { code } => {
if !r_is_at_prompt().load(Ordering::Acquire) {
let _ = reply.send(IpcResponse::error(
R_NOT_AT_PROMPT,
"R is not at the command prompt".to_string(),
));
return;
}
*pending_ipc_operation()
.lock()
.unwrap_or_else(|e| e.into_inner()) = Some(PendingIpcOperation {
kind: PendingIpcKind::UserInput { reply },
code,
});
fire_break_signal();
}
IpcMethod::Session => unreachable!("Session handled above"),
}
}
fn fire_break_signal() {
if let Some(signal) = BREAK_SIGNAL.get() {
signal.store(true, Ordering::Relaxed);
}
}
pub fn take_pending_ipc_operation() -> Option<PendingIpcOperation> {
pending_ipc_operation()
.lock()
.unwrap_or_else(|e| e.into_inner())
.take()
}
pub fn reject_operation_user_typing(op: PendingIpcOperation, buffer: &str) {
const MAX_BUFFER_CHARS: usize = 1024;
let original_len = buffer.chars().count();
let truncated = original_len > MAX_BUFFER_CHARS;
let preview: String = buffer.chars().take(MAX_BUFFER_CHARS).collect();
let response = IpcResponse::error_with_data(
USER_IS_TYPING,
"User is typing in the console".to_string(),
serde_json::json!({
"buffer": preview,
"buffer_truncated": truncated,
"buffer_original_length": original_len,
}),
);
match op.kind {
PendingIpcKind::SilentEvaluate { reply }
| PendingIpcKind::VisibleEvaluate { reply, .. }
| PendingIpcKind::UserInput { reply } => {
let _ = reply.send(response);
}
}
}
pub fn setup_visible_eval(
reply: tokio::sync::oneshot::Sender<IpcResponse>,
timeout: std::time::Duration,
) {
r_is_at_prompt().store(false, Ordering::Release);
arf_libr::start_ipc_capture(true);
*pending_visible_eval()
.lock()
.unwrap_or_else(|e| e.into_inner()) = Some(PendingVisibleEval {
reply,
deadline: std::time::Instant::now()
.checked_add(timeout)
.unwrap_or_else(|| std::time::Instant::now() + DEFAULT_EVAL_TIMEOUT),
timeout,
});
}
pub fn run_silent_eval(code: &str, reply: tokio::sync::oneshot::Sender<IpcResponse>) {
r_is_at_prompt().store(false, Ordering::Release);
let result = capture::evaluate_with_capture(code, false);
r_is_at_prompt().store(true, Ordering::Release);
let _ = reply.send(IpcResponse::Evaluate(result));
}
pub fn accept_user_input(reply: tokio::sync::oneshot::Sender<IpcResponse>) {
let _ = reply.send(IpcResponse::UserInput(UserInputResult { accepted: true }));
}
pub fn set_r_at_prompt(at_prompt: bool) {
r_is_at_prompt().store(at_prompt, Ordering::Release);
}
pub fn set_headless_shutdown(flag: Arc<AtomicBool>) {
let _ = HEADLESS_SHUTDOWN.set(flag);
}
pub fn trigger_headless_shutdown() -> bool {
if let Some(flag) = HEADLESS_SHUTDOWN.get() {
flag.store(true, Ordering::Release);
true
} else {
false
}
}
fn arf_session_base(meta: &SessionMeta) -> SessionResult {
SessionResult {
arf_version: env!("CARGO_PKG_VERSION").to_string(),
pid: std::process::id(),
os: std::env::consts::OS.to_string(),
arch: std::env::consts::ARCH.to_string(),
socket_path: meta.socket_path.clone(),
started_at: meta.started_at.clone(),
log_file: meta.log_file.clone(),
history_session_id: meta.history_session_id,
r: None,
r_unavailable_reason: None,
hint: None,
}
}
#[derive(Clone)]
struct SessionMeta {
socket_path: String,
started_at: String,
log_file: Option<String>,
history_session_id: Option<i64>,
}
static SESSION_META: OnceLock<Mutex<SessionMeta>> = OnceLock::new();
pub fn clear_history_session_id() {
if let Some(m) = SESSION_META.get() {
let mut meta = m.lock().unwrap_or_else(|e| e.into_inner());
if meta.history_session_id.take().is_some() {
session::clear_session_history_id(std::process::id());
}
}
}
pub fn set_headless_history(history: reedline::SqliteBackedHistory) {
if HEADLESS_HISTORY.set(Mutex::new(history)).is_err() {
log::warn!(
"Headless history backend already initialized; ignoring duplicate set_headless_history call"
);
}
}
pub fn set_history_db_info(path: PathBuf, session_id: Option<reedline::HistorySessionId>) {
if HISTORY_DB_INFO.set((path, session_id)).is_err() {
log::warn!(
"History DB info already initialized; ignoring duplicate set_history_db_info call"
);
}
}
pub(crate) enum HistoryQueryError {
InvalidParams(String),
Internal(String),
}
pub(crate) fn query_history(params: &HistoryParams) -> Result<HistoryResult, HistoryQueryError> {
if params.limit < 1 {
return Err(HistoryQueryError::InvalidParams(format!(
"limit must be positive, got {}",
params.limit
)));
}
let (db_path, session_id) = HISTORY_DB_INFO.get().ok_or_else(|| {
HistoryQueryError::Internal(
"History is not available (no history database configured)".to_string(),
)
})?;
if !params.all_sessions && session_id.is_none() {
return Err(HistoryQueryError::Internal(
"Session ID is not available; use all_sessions=true to query without session scope"
.to_string(),
));
}
let since_ms = if let Some(ref since) = params.since {
Some(
chrono::DateTime::parse_from_rfc3339(since)
.map(|dt| dt.with_timezone(&chrono::Utc))
.or_else(|_| {
chrono::NaiveDate::parse_from_str(since, "%Y-%m-%d")
.map(|d| d.and_hms_opt(0, 0, 0).unwrap().and_utc())
})
.map_err(|_| {
HistoryQueryError::InvalidParams(format!(
"Invalid 'since' format: {since}. \
Use RFC 3339 (e.g. '2026-03-29T00:00:00Z') or date (e.g. '2026-03-29')"
))
})?
.timestamp_millis(),
)
} else {
None
};
let db =
rusqlite::Connection::open_with_flags(db_path, rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY)
.map_err(|e| {
HistoryQueryError::Internal(format!("Failed to open history database: {e}"))
})?;
let mut sql = String::from(
"SELECT command_line, start_timestamp, session_id, cwd, exit_status \
FROM history WHERE 1=1",
);
let mut params_vec: Vec<Box<dyn rusqlite::types::ToSql>> = Vec::new();
if !params.all_sessions
&& let Some(sid) = session_id
{
sql.push_str(" AND session_id = ?");
params_vec.push(Box::new(i64::from(*sid)));
}
if let Some(ref cwd) = params.cwd {
sql.push_str(" AND cwd = ?");
params_vec.push(Box::new(cwd.clone()));
}
if let Some(ref grep) = params.grep {
sql.push_str(" AND instr(command_line, ?) >= 1");
params_vec.push(Box::new(grep.clone()));
}
if let Some(ms) = since_ms {
sql.push_str(" AND start_timestamp >= ?");
params_vec.push(Box::new(ms));
}
sql.push_str(" ORDER BY id DESC LIMIT ?");
params_vec.push(Box::new(params.limit));
let param_refs: Vec<&dyn rusqlite::types::ToSql> = params_vec.iter().map(|p| &**p).collect();
let mut stmt = db
.prepare(&sql)
.map_err(|e| HistoryQueryError::Internal(format!("Failed to prepare query: {e}")))?;
let entries = stmt
.query_map(param_refs.as_slice(), |row| {
let ts_ms: Option<i64> = row.get(1)?;
let timestamp: Option<String> = ts_ms.and_then(|ms| {
chrono::Utc
.timestamp_millis_opt(ms)
.single()
.map(|t| t.to_rfc3339())
});
let sid: Option<i64> = row.get(2)?;
Ok(HistoryEntry {
command: row.get(0)?,
timestamp,
cwd: row.get(3)?,
exit_status: row.get(4)?,
session_id: sid,
})
})
.map_err(|e| HistoryQueryError::Internal(format!("History query failed: {e}")))?
.collect::<Result<Vec<_>, _>>()
.map_err(|e| HistoryQueryError::Internal(format!("Failed to read history row: {e}")))?;
Ok(HistoryResult {
entries,
session_id: session_id.map(i64::from),
})
}
fn save_to_headless_history(code: &str, exit_status: Option<i64>) {
let Some(h) = HEADLESS_HISTORY.get() else {
return;
};
let Ok(mut history) = h.lock() else {
log::warn!("Headless history lock poisoned, skipping save");
return;
};
use reedline::History;
let mut item = reedline::HistoryItem::from_command_line(code);
item.start_timestamp = Some(chrono::Utc::now());
item.hostname = Some(gethostname::gethostname().to_string_lossy().into_owned());
item.cwd = std::env::current_dir()
.ok()
.map(|p| p.to_string_lossy().into_owned());
item.exit_status = exit_status;
item.session_id = HISTORY_DB_INFO.get().and_then(|(_, sid)| *sid);
if let Err(e) = history.save(item) {
log::warn!("Failed to save headless history: {}", e);
}
}
pub(in crate::ipc) fn set_session_meta(
socket_path: String,
started_at: String,
log_file: Option<String>,
history_session_id: Option<i64>,
) {
let meta = SessionMeta {
socket_path,
started_at,
log_file,
history_session_id,
};
match SESSION_META.get() {
Some(m) => *m.lock().unwrap_or_else(|e| e.into_inner()) = meta,
None => {
let _ = SESSION_META.set(Mutex::new(meta));
}
}
}
fn current_session_meta() -> SessionMeta {
match SESSION_META.get() {
Some(m) => m.lock().unwrap_or_else(|e| e.into_inner()).clone(),
None => SessionMeta {
socket_path: "<uninitialized_socket_path>".to_string(),
started_at: "<uninitialized_started_at>".to_string(),
log_file: None,
history_session_id: None,
},
}
}
pub(in crate::ipc) fn collect_session_result(try_r: bool, reason: &str) -> SessionResult {
let meta = current_session_meta();
let mut result = arf_session_base(&meta);
if !try_r || !r_is_at_prompt().load(Ordering::Acquire) {
let reason = if reason.is_empty() {
"R is busy evaluating another expression"
} else {
reason
};
result.r_unavailable_reason = Some(reason.to_string());
result.hint = Some(if reason.contains("alternate mode") {
"Exit the current mode (shell, browser) to make R session info available.".to_string()
} else if reason.contains("pending") {
"Wait for the current IPC operation to complete, then retry.".to_string()
} else if reason.contains("Main thread") || reason.contains("handler dropped") {
"The arf process may be shutting down or unresponsive.".to_string()
} else if reason.contains("Timed out") {
"R may be busy with a long-running operation. Retry later.".to_string()
} else {
"R session information will be available when R returns to the prompt. \
Retry 'arf ipc session' later, or use 'arf ipc eval' with a timeout to wait."
.to_string()
});
return result;
}
match collect_r_session_info() {
Some(r_info) => {
result.r = Some(r_info);
}
None => {
result.r_unavailable_reason =
Some("Failed to collect R session information".to_string());
result.hint = Some(
"R may not be fully initialized. Try again later or use \
'arf ipc eval \"sessionInfo()\"' for raw output."
.to_string(),
);
}
}
result
}
fn collect_r_session_info() -> Option<RSessionInfo> {
let version = eval_r_scalar(r#"invisible(paste0(R.version$major, ".", R.version$minor))"#)
.unwrap_or_default();
if version.is_empty() {
return None;
}
Some(RSessionInfo {
version,
platform: eval_r_scalar("invisible(R.version$platform)").unwrap_or_default(),
locale: eval_r_scalar("invisible(Sys.getlocale())").unwrap_or_default(),
cwd: eval_r_scalar("invisible(getwd())").unwrap_or_default(),
loaded_namespaces: eval_r_character_vector("invisible(loadedNamespaces())")
.unwrap_or_default(),
attached_packages: eval_r_character_vector("invisible(.packages())").unwrap_or_default(),
lib_paths: eval_r_character_vector("invisible(.libPaths())").unwrap_or_default(),
})
}
fn eval_r_scalar(code: &str) -> Option<String> {
match arf_harp::eval_string(code) {
Ok(robj) => extract_r_string(robj.sexp()),
Err(e) => {
log::debug!("eval_r_scalar failed for `{code}`: {e}");
None
}
}
}
fn eval_r_character_vector(code: &str) -> Option<Vec<String>> {
match arf_harp::eval_string(code) {
Ok(robj) => extract_r_strings(robj.sexp()),
Err(e) => {
log::debug!("eval_r_character_vector failed for `{code}`: {e}");
None
}
}
}
fn extract_r_string(sexp: arf_libr::SEXP) -> Option<String> {
let lib = arf_libr::r_library().ok()?;
unsafe {
if (lib.rf_isstring)(sexp) == 0 || (lib.rf_length)(sexp) == 0 {
return None;
}
let elt = (lib.string_elt)(sexp, 0);
let cstr = (lib.r_charsxp)(elt);
if cstr.is_null() {
return None;
}
std::ffi::CStr::from_ptr(cstr)
.to_str()
.ok()
.map(|s| s.to_string())
}
}
fn extract_r_strings(sexp: arf_libr::SEXP) -> Option<Vec<String>> {
let lib = arf_libr::r_library().ok()?;
unsafe {
if (lib.rf_isstring)(sexp) == 0 {
return None;
}
let len = (lib.rf_length)(sexp) as isize;
let mut result = Vec::with_capacity(len as usize);
for i in 0..len {
let elt = (lib.string_elt)(sexp, i);
let cstr = (lib.r_charsxp)(elt);
if cstr.is_null() {
result.push(String::new());
} else if let Ok(s) = std::ffi::CStr::from_ptr(cstr).to_str() {
result.push(s.to_string());
} else {
result.push(String::new());
}
}
Some(result)
}
}
pub fn headless_poll_and_process() -> bool {
let receiver = match IPC_RECEIVER.get() {
Some(r) => r,
None => return false,
};
let rx = match receiver.try_lock() {
Ok(rx) => rx,
Err(_) => return false,
};
let rx = match rx.as_ref() {
Some(rx) => rx,
None => return false,
};
let mut processed = false;
while let Ok(request) = rx.try_recv() {
processed = true;
headless_handle_request(request);
}
processed
}
fn headless_handle_request(request: IpcRequest) {
let IpcRequest { method, reply } = request;
match method {
IpcMethod::Evaluate { code, visible, .. } => {
r_is_at_prompt().store(false, Ordering::Release);
let result = capture::evaluate_with_capture(&code, visible);
r_is_at_prompt().store(true, Ordering::Release);
let has_error = result.error.is_some();
let _ = reply.send(IpcResponse::Evaluate(result));
if !code.trim().is_empty() {
let exit_status = if has_error { 1 } else { 0 };
save_to_headless_history(&code, Some(exit_status));
}
}
IpcMethod::UserInput { code } => {
r_is_at_prompt().store(false, Ordering::Release);
let eval_result = arf_harp::eval_string(&code);
r_is_at_prompt().store(true, Ordering::Release);
let exit_status;
match eval_result {
Ok(_) => {
exit_status = 0;
let _ = reply.send(IpcResponse::UserInput(UserInputResult { accepted: true }));
}
Err(e) => {
exit_status = 1;
log::warn!("Headless user_input evaluation error: {}", e);
let _ = reply.send(IpcResponse::UserInput(UserInputResult { accepted: false }));
}
}
if !code.trim().is_empty() {
save_to_headless_history(&code, Some(exit_status));
}
}
IpcMethod::Session => {
let result = collect_session_result(true, "");
let _ = reply.send(IpcResponse::Session(Box::new(result)));
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use serial_test::serial;
struct GlobalStateGuard;
impl Drop for GlobalStateGuard {
fn drop(&mut self) {
set_in_alternate_mode(false);
set_r_at_prompt(false);
}
}
#[test]
#[serial]
fn test_alternate_mode_flag_and_request_rejection() {
set_in_alternate_mode(false);
set_r_at_prompt(false);
let _guard = GlobalStateGuard;
assert!(!is_in_alternate_mode());
set_in_alternate_mode(true);
assert!(is_in_alternate_mode());
set_in_alternate_mode(false);
assert!(!is_in_alternate_mode());
set_r_at_prompt(true);
set_in_alternate_mode(true);
{
let (reply_tx, reply_rx) = tokio::sync::oneshot::channel();
let request = IpcRequest {
method: IpcMethod::UserInput {
code: "1+1".to_string(),
},
reply: reply_tx,
};
handle_request(request);
match reply_rx.blocking_recv().unwrap() {
IpcResponse::Error { code, .. } => assert_eq!(code, R_NOT_AT_PROMPT),
_ => panic!("Expected R_NOT_AT_PROMPT error for user_input"),
}
}
{
let (reply_tx, reply_rx) = tokio::sync::oneshot::channel();
let request = IpcRequest {
method: IpcMethod::Evaluate {
code: "1+1".to_string(),
visible: false,
timeout_ms: None,
},
reply: reply_tx,
};
handle_request(request);
match reply_rx.blocking_recv().unwrap() {
IpcResponse::Error { code, .. } => assert_eq!(code, R_NOT_AT_PROMPT),
_ => panic!("Expected R_NOT_AT_PROMPT error for evaluate"),
}
}
}
#[test]
#[serial]
fn test_session_returns_arf_only_in_various_states() {
set_in_alternate_mode(false);
set_r_at_prompt(false);
let _guard = GlobalStateGuard;
fn send_session() -> protocol::SessionResult {
let (reply_tx, reply_rx) = tokio::sync::oneshot::channel();
let request = IpcRequest {
method: IpcMethod::Session,
reply: reply_tx,
};
handle_request(request);
match reply_rx.blocking_recv().unwrap() {
IpcResponse::Session(result) => *result,
_ => panic!("Expected Session response"),
}
}
set_in_alternate_mode(true);
set_r_at_prompt(true);
{
let result = send_session();
assert!(result.r.is_none());
let reason = result.r_unavailable_reason.unwrap();
assert!(
reason.contains("alternate mode"),
"Expected alternate mode reason, got: {reason}"
);
}
set_in_alternate_mode(false);
set_r_at_prompt(false);
{
let result = send_session();
assert!(result.r.is_none());
assert!(result.r_unavailable_reason.is_some());
}
set_r_at_prompt(true);
{
let (dummy_tx, _dummy_rx) = tokio::sync::oneshot::channel();
*pending_ipc_operation()
.lock()
.unwrap_or_else(|e| e.into_inner()) = Some(PendingIpcOperation {
kind: PendingIpcKind::SilentEvaluate { reply: dummy_tx },
code: "dummy".to_string(),
});
let result = send_session();
assert!(result.r.is_none());
let reason = result.r_unavailable_reason.unwrap();
assert!(
reason.contains("pending"),
"Expected pending reason, got: {reason}"
);
let _ = take_pending_ipc_operation();
}
}
}