use crate::ipc::protocol::JsonRpcResponse;
use crate::ipc::session::{find_session, list_sessions};
use crate::output::{self, write_json};
use anyhow::{Context, Result};
const DEFAULT_TRANSPORT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(300);
const EXIT_CLIENT: i32 = 2;
const EXIT_SESSION: i32 = 3;
const EXIT_PROTOCOL: i32 = 4;
fn print_json_or_exit(value: &serde_json::Value) {
if let Err(error) = output::print_json(value) {
exit_error(
EXIT_CLIENT,
"OUTPUT_ERROR",
&format!("Failed to write JSON output: {error:#}"),
None,
None,
);
}
}
fn exit_error(
exit_code: i32,
code: &str,
message: &str,
hint: Option<&str>,
data: Option<&serde_json::Value>,
) -> ! {
let error = serde_json::json!({
"error": {
"code": code,
"message": message,
"hint": hint,
"data": data,
}
});
let pretty = std::io::IsTerminal::is_terminal(&std::io::stderr());
let write_result = {
let mut stderr = std::io::stderr().lock();
write_json(&mut stderr, &error, pretty).and_then(|_| {
use std::io::Write;
writeln!(stderr).context("Failed to write JSON newline")
})
};
if let Err(write_error) = write_result {
eprintln!("Error writing JSON error: {write_error:#}");
}
std::process::exit(exit_code);
}
fn rpc_error_info(code: i32) -> (&'static str, Option<&'static str>) {
use crate::ipc::protocol::*;
match code {
R_BUSY => (
"R_BUSY",
Some(
"R is executing code. Wait for it to finish, or use \
'arf ipc session' to check status.",
),
),
R_NOT_AT_PROMPT => (
"R_NOT_AT_PROMPT",
Some(
"R is in browser/debug mode or a menu. If triggered by \
browser(), use 'arf ipc send \"Q\"' to exit the debugger.",
),
),
INPUT_ALREADY_PENDING => (
"INPUT_ALREADY_PENDING",
Some(
"Another IPC input is already queued. Wait for it to \
be processed before sending more.",
),
),
USER_IS_TYPING => (
"USER_IS_TYPING",
Some(
"The user is typing in the REPL. Check error.data.buffer \
for the current input. Wait for them to finish or clear \
their input.",
),
),
INCOMPLETE_INPUT => (
"INCOMPLETE_INPUT",
Some("Complete the R expression before sending it over IPC."),
),
R_EVAL_NOT_ALLOWED => (
"R_EVAL_NOT_ALLOWED",
Some(
"--visible needs no allowlist entry, because it runs the code where the \
session shows it rather than silently. The allowlist itself is set only at \
startup, through [ipc.eval].allowed_functions or repeated \
--ipc-eval-allow-function options, and lifted with --ipc-eval-unrestricted.",
),
),
INPUT_NOT_APPROVED => (
"INPUT_NOT_APPROVED",
Some("Approve the request in the interactive REPL, then send it again."),
),
PARSE_ERROR => ("PARSE_ERROR", None),
INVALID_REQUEST => ("INVALID_REQUEST", None),
METHOD_NOT_FOUND => ("METHOD_NOT_FOUND", None),
INVALID_PARAMS => ("INVALID_PARAMS", None),
INTERNAL_ERROR => ("INTERNAL_ERROR", None),
_ => ("PROTOCOL_ERROR", None),
}
}
fn handle_response(response: JsonRpcResponse) {
if let Some(ref error) = response.error {
let (code, hint) = rpc_error_info(error.code);
exit_error(
EXIT_PROTOCOL,
code,
&error.message,
hint,
error.data.as_ref(),
);
}
match response.result {
Some(result) => print_json_or_exit(&result),
None => {
exit_error(
EXIT_PROTOCOL,
"EMPTY_RESPONSE",
"Server returned a response with neither result nor error (possible server bug)",
None,
None,
);
}
}
}
pub fn cmd_list() {
let sessions = list_sessions();
let sessions_json: Vec<serde_json::Value> = sessions
.iter()
.map(|s| {
serde_json::to_value(s).unwrap_or_else(|e| {
exit_error(
EXIT_PROTOCOL,
"SERIALIZATION_ERROR",
&format!("Failed to serialize session info: {e}"),
Some("This is likely a bug in arf."),
None,
);
})
})
.collect();
print_json_or_exit(&serde_json::json!({ "sessions": sessions_json }));
}
fn resolve_session(pid: Option<u32>) -> crate::ipc::session::SessionInfo {
match find_session(pid) {
Some(session) => session,
None => {
if let Some(p) = pid {
exit_error(
EXIT_SESSION,
"SESSION_NOT_FOUND",
&format!("No active arf session with PID {p}"),
Some("Use 'arf ipc list' to see active sessions."),
None,
);
} else {
let sessions = list_sessions();
if sessions.is_empty() {
exit_error(
EXIT_SESSION,
"SESSION_NOT_FOUND",
"No active arf sessions found",
Some("Start arf with --with-ipc to enable IPC."),
None,
);
} else {
exit_error(
EXIT_SESSION,
"SESSION_AMBIGUOUS",
"Multiple arf sessions running",
Some(
"Specify --pid to select one. Use 'arf ipc list' \
to see active sessions.",
),
None,
);
}
}
}
}
}
fn require_stdin_not_tty() {
use std::io::IsTerminal;
if std::io::stdin().is_terminal() {
exit_error(
EXIT_CLIENT,
"NO_CODE_PROVIDED",
"No code provided and stdin is a terminal",
Some("Pass code as an argument or pipe it via stdin."),
None,
);
}
}
fn read_stdin_code() -> String {
use std::io::Read;
let mut buf = String::new();
std::io::stdin()
.read_to_string(&mut buf)
.unwrap_or_else(|e| {
exit_error(
EXIT_CLIENT,
"STDIN_READ_ERROR",
&format!("Failed to read code from stdin: {e}"),
None,
None,
);
});
buf
}
pub fn cmd_eval(code: Option<&str>, pid: Option<u32>, visible: bool, timeout_ms: Option<u64>) {
if code.is_none() {
require_stdin_not_tty();
}
let session = resolve_session(pid);
let owned;
let code = match code {
Some(c) => c,
None => {
owned = read_stdin_code();
&owned
}
};
let mut params = serde_json::json!({ "code": code, "visible": visible });
if let Some(ms) = timeout_ms {
params["timeout_ms"] = serde_json::json!(ms);
}
let request = serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "evaluate",
"params": params
});
let transport_timeout = match timeout_ms {
Some(ms) => std::time::Duration::from_millis(ms.saturating_add(5000)),
None => DEFAULT_TRANSPORT_TIMEOUT + std::time::Duration::from_secs(5),
};
let response = send_request(&session.socket_path, &request, transport_timeout);
handle_response(response);
}
pub fn cmd_send(code: Option<&str>, pid: Option<u32>) {
if code.is_none() {
require_stdin_not_tty();
}
let session = resolve_session(pid);
let owned;
let code = match code {
Some(c) => c,
None => {
owned = read_stdin_code();
&owned
}
};
let request = serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "user_input",
"params": { "code": code }
});
let response = send_request(&session.socket_path, &request, DEFAULT_TRANSPORT_TIMEOUT);
handle_response(response);
}
pub fn cmd_shutdown(pid: Option<u32>) {
let session = resolve_session(pid);
let request = serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "shutdown",
"params": {}
});
let response = send_request(&session.socket_path, &request, DEFAULT_TRANSPORT_TIMEOUT);
handle_response(response);
}
pub fn cmd_session(pid: Option<u32>) {
let session = resolve_session(pid);
let request = serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "session",
"params": {}
});
let transport_timeout = std::time::Duration::from_secs(15);
let response = send_request(&session.socket_path, &request, transport_timeout);
handle_response(response);
}
pub fn cmd_history(
pid: Option<u32>,
limit: i64,
all_sessions: bool,
cwd: Option<&str>,
grep: Option<&str>,
since: Option<&str>,
) {
let session = resolve_session(pid);
let mut params = serde_json::json!({
"limit": limit,
"all_sessions": all_sessions,
});
if let Some(cwd) = cwd {
params["cwd"] = serde_json::Value::String(cwd.to_string());
}
if let Some(grep) = grep {
params["grep"] = serde_json::Value::String(grep.to_string());
}
if let Some(since) = since {
params["since"] = serde_json::Value::String(since.to_string());
}
let request = serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "history",
"params": params
});
let transport_timeout = std::time::Duration::from_secs(15);
let response = send_request(&session.socket_path, &request, transport_timeout);
handle_response(response);
}
fn send_request(
socket_path: &str,
request: &serde_json::Value,
timeout: std::time::Duration,
) -> JsonRpcResponse {
match send_request_inner(socket_path, request, timeout) {
Ok(response) => response,
Err(e) => {
let is_protocol = e
.chain()
.any(|cause| cause.downcast_ref::<serde_json::Error>().is_some());
if is_protocol {
exit_error(
EXIT_PROTOCOL,
"PROTOCOL_ERROR",
&format!("{e:#}"),
Some("Received an invalid or malformed response from the arf session."),
None,
);
} else {
exit_error(
EXIT_CLIENT,
"TRANSPORT_ERROR",
&format!("{e:#}"),
Some("Check that the arf session is running and IPC is enabled."),
None,
);
}
}
}
}
fn send_request_inner(
socket_path: &str,
request: &serde_json::Value,
timeout: std::time::Duration,
) -> Result<JsonRpcResponse> {
let body = serde_json::to_string(request)?;
#[cfg(unix)]
{
use std::io::{Read, Write};
use std::os::unix::net::UnixStream;
let http_request = format!(
"POST / HTTP/1.1\r\n\
Host: localhost\r\n\
Content-Type: application/json\r\n\
Content-Length: {}\r\n\
Connection: close\r\n\
\r\n{}",
body.len(),
body
);
let mut stream = UnixStream::connect(socket_path)
.with_context(|| format!("Failed to connect to {socket_path}"))?;
stream.set_read_timeout(Some(timeout))?;
stream.write_all(http_request.as_bytes())?;
stream.shutdown(std::net::Shutdown::Write)?;
let mut response_buf = Vec::new();
stream.read_to_end(&mut response_buf)?;
parse_http_response(&response_buf)
}
#[cfg(windows)]
{
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::windows::named_pipe::ClientOptions;
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.context("Failed to create tokio runtime")?;
rt.block_on(async {
let mut pipe = ClientOptions::new()
.open(socket_path)
.with_context(|| format!("Failed to connect to {socket_path}"))?;
pipe.write_all(body.as_bytes()).await?;
pipe.flush().await?;
let mut response_buf = Vec::new();
match tokio::time::timeout(timeout, pipe.read_to_end(&mut response_buf)).await {
Ok(result) => result?,
Err(_) => anyhow::bail!("Request timed out after {}s", timeout.as_secs()),
};
parse_http_response(&response_buf)
})
}
}
fn parse_http_response(data: &[u8]) -> Result<JsonRpcResponse> {
let text = String::from_utf8_lossy(data);
let body = if let Some(pos) = text.find("\r\n\r\n") {
&text[pos + 4..]
} else {
&text
};
serde_json::from_str(body).context("Failed to parse JSON-RPC response")
}