use anyhow::{Context, Result, bail};
use chrono::{DateTime, NaiveDateTime, Utc};
use reedline::{HistoryItem, HistoryItemId, HistorySessionId};
use std::collections::{HashMap, HashSet};
use std::fs::File;
use std::io::{BufRead, BufReader};
use std::path::{Path, PathBuf};
use std::time::Duration;
use super::metadata::HistoryExtraInfo;
use super::store::HistoryStore;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ImportMode {
R,
Shell,
Browse,
Unspecified,
Unsupported(String),
}
impl ImportMode {
fn from_external(mode: Option<&str>) -> Self {
match mode {
Some("r") => Self::R,
Some("shell") => Self::Shell,
Some("browse") => Self::Browse,
Some(mode) => Self::Unsupported(mode.to_owned()),
None => Self::Unspecified,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ImportEntry {
pub mode: ImportMode,
pub item: HistoryItem<HistoryExtraInfo>,
}
impl ImportEntry {
pub fn new(command: impl Into<String>) -> Self {
Self {
mode: ImportMode::Unspecified,
item: HistoryItem {
id: None,
start_timestamp: None,
command_line: command.into(),
session_id: None,
hostname: None,
cwd: None,
duration: None,
exit_status: None,
more_info: None,
},
}
}
pub fn with_mode(mut self, mode: ImportMode) -> Self {
self.mode = mode;
self
}
}
#[derive(Debug, Default)]
pub struct ParsedImport {
pub entries: Vec<ImportEntry>,
pub warnings: Vec<String>,
}
impl std::ops::Deref for ParsedImport {
type Target = [ImportEntry];
fn deref(&self) -> &Self::Target {
&self.entries
}
}
#[derive(Debug, Default, PartialEq, Eq)]
pub struct ImportResult {
pub r_imported: usize,
pub shell_imported: usize,
pub skipped: usize,
pub duplicates_skipped: usize,
pub duplicates_repaired: usize,
pub warnings: Vec<String>,
}
impl ImportResult {
#[allow(dead_code)]
pub fn total_imported(&self) -> usize {
self.r_imported + self.shell_imported
}
}
pub fn default_radian_path() -> PathBuf {
dirs::home_dir()
.map(|h| h.join(".radian_history"))
.unwrap_or_else(|| PathBuf::from(".radian_history"))
}
pub fn default_r_history_path() -> PathBuf {
if let Ok(path) = std::env::var("R_HISTFILE") {
return PathBuf::from(path);
}
PathBuf::from(".Rhistory")
}
pub fn parse_radian_history(path: &Path) -> Result<ParsedImport> {
let file = File::open(path)
.with_context(|| format!("Failed to open radian history: {}", path.display()))?;
let reader = BufReader::new(file);
let mut entries = Vec::new();
let mut current_timestamp: Option<DateTime<Utc>> = None;
let mut current_mode: Option<String> = None;
let mut current_lines: Vec<String> = Vec::new();
for line_result in reader.lines() {
let line = line_result.with_context(|| "Failed to read line from radian history")?;
if line.starts_with("# time: ") {
if !current_lines.is_empty() {
let command = current_lines.join("\n");
let mut entry = ImportEntry::new(command)
.with_mode(ImportMode::from_external(current_mode.take().as_deref()));
entry.item.start_timestamp = current_timestamp;
entries.push(entry);
current_lines.clear();
}
current_mode = None;
let time_str = line.trim_start_matches("# time: ").trim();
let time_str = time_str.trim_end_matches(" UTC");
current_timestamp = NaiveDateTime::parse_from_str(time_str, "%Y-%m-%d %H:%M:%S")
.ok()
.map(|naive| naive.and_utc());
} else if line.starts_with("# mode: ") {
current_mode = Some(line.trim_start_matches("# mode: ").trim().to_string());
} else if let Some(content) = line.strip_prefix('+') {
let content = content.strip_suffix('\r').unwrap_or(content);
current_lines.push(content.to_string());
} else if line.trim().is_empty() {
if !current_lines.is_empty() {
let command = current_lines.join("\n");
let mut entry = ImportEntry::new(command)
.with_mode(ImportMode::from_external(current_mode.take().as_deref()));
entry.item.start_timestamp = current_timestamp;
entries.push(entry);
current_lines.clear();
current_timestamp = None;
}
}
}
if !current_lines.is_empty() {
let command = current_lines.join("\n");
let mut entry = ImportEntry::new(command)
.with_mode(ImportMode::from_external(current_mode.take().as_deref()));
entry.item.start_timestamp = current_timestamp;
entries.push(entry);
}
Ok(ParsedImport {
entries,
warnings: Vec::new(),
})
}
pub fn parse_r_history(path: &Path) -> Result<ParsedImport> {
let file = File::open(path)
.with_context(|| format!("Failed to open R history: {}", path.display()))?;
let reader = BufReader::new(file);
let mut entries = Vec::new();
for line_result in reader.lines() {
let line = line_result.with_context(|| "Failed to read line from R history")?;
let content = line.trim_end();
if !content.trim().is_empty() {
entries.push(ImportEntry::new(content.to_string()).with_mode(ImportMode::R));
}
}
Ok(ParsedImport {
entries,
warnings: Vec::new(),
})
}
pub fn parse_arf_history(path: &Path) -> Result<ParsedImport> {
if !path.exists() {
bail!("arf history database not found: {}", path.display());
}
let is_shell = path
.file_name()
.and_then(|n| n.to_str())
.is_some_and(|n| n == "shell.db");
let mode = if is_shell {
ImportMode::Shell
} else {
ImportMode::R
};
let db =
rusqlite::Connection::open_with_flags(path, rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY)
.with_context(|| format!("Failed to open arf history database: {}", path.display()))?;
if !table_exists(&db, "history")? {
bail!(
"File '{}' does not look like an arf history database: missing history table",
path.display()
);
}
read_history_table(&db, path, "history", mode).with_context(|| {
format!(
"File '{}' does not look like an arf history database",
path.display()
)
})
}
pub struct ImportTargets {
pub r_history: HistoryStore,
pub shell_history: HistoryStore,
}
#[derive(Clone)]
pub struct DedupSet {
rows: Vec<DedupRow>,
command_timestamps_by_row: HashMap<(String, i64), Vec<usize>>,
command_rows: HashMap<String, Vec<usize>>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum MetadataState {
Null,
Valid,
Malformed,
}
#[derive(Debug, Clone)]
struct DedupRow {
id: HistoryItemId,
has_session_id: bool,
has_hostname: bool,
has_cwd: bool,
has_duration: bool,
has_exit_status: bool,
metadata: MetadataState,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum DuplicateAction {
NotDuplicate,
Skip,
Repair(HistoryItemId),
Ambiguous,
Malformed,
}
impl DedupSet {
pub fn from_history(history: &HistoryStore) -> Result<Self> {
let path = history
.path()
.ok_or_else(|| anyhow::anyhow!("history store has no persistent path"))?;
Self::from_connection(
rusqlite::Connection::open(path)
.with_context(|| format!("Failed to read history database: {}", path.display()))?,
)
}
pub fn from_db(path: &Path) -> Result<Self> {
use rusqlite::{Connection, OpenFlags};
let db = Connection::open_with_flags(path, OpenFlags::SQLITE_OPEN_READ_ONLY)
.with_context(|| format!("Failed to open history database: {}", path.display()))?;
Self::from_connection(db)
}
fn from_connection(db: rusqlite::Connection) -> Result<Self> {
let columns = HistoryTableColumns::read(&db, "history")?;
let query = format!(
"SELECT id, command_line, start_timestamp, {}, {}, {}, {}, {}, {} FROM history",
columns.expression("session_id"),
columns.expression("hostname"),
columns.expression("cwd"),
columns.expression("duration_ms"),
columns.expression("exit_status"),
columns.expression("more_info"),
);
let mut stmt = db
.prepare(&query)
.context("Failed to query history for dedup")?;
let rows = stmt
.query_map([], |row| {
Ok((
HistoryItemId::new(row.get::<_, i64>(0)?),
row.get::<_, String>(1)?,
row.get::<_, Option<i64>>(2)?,
row.get::<_, Option<i64>>(3)?.is_some(),
row.get::<_, Option<String>>(4)?.is_some(),
row.get::<_, Option<String>>(5)?.is_some(),
row.get::<_, Option<i64>>(6)?.is_some(),
row.get::<_, Option<i64>>(7)?.is_some(),
row.get::<_, Option<String>>(8)?,
))
})
.context("Failed to query history for dedup")?;
let mut set = Self {
rows: Vec::new(),
command_timestamps_by_row: HashMap::new(),
command_rows: HashMap::new(),
};
for row in rows {
let (
id,
command,
timestamp_millis,
has_session_id,
has_hostname,
has_cwd,
has_duration,
has_exit_status,
raw_metadata,
) = row.context("Failed to read history row")?;
let metadata = match raw_metadata.as_deref() {
None => MetadataState::Null,
Some(raw) => match serde_json::from_str::<HistoryExtraInfo>(raw) {
Ok(_) => MetadataState::Valid,
Err(_) => MetadataState::Malformed,
},
};
let row_index = set.rows.len();
set.command_rows
.entry(command.clone())
.or_default()
.push(row_index);
if let Some(ms) = timestamp_millis {
set.command_timestamps_by_row
.entry((command.clone(), ms))
.or_default()
.push(row_index);
}
set.rows.push(DedupRow {
id,
has_session_id,
has_hostname,
has_cwd,
has_duration,
has_exit_status,
metadata,
});
}
Ok(set)
}
fn duplicate_action(
&self,
command: &str,
timestamp: Option<&DateTime<Utc>>,
item: &HistoryItem<HistoryExtraInfo>,
) -> DuplicateAction {
let row_indices = if let Some(ts) = timestamp {
self.command_timestamps_by_row
.get(&(command.to_string(), ts.timestamp_millis()))
} else {
self.command_rows.get(command)
};
let Some(row_indices) = row_indices else {
return DuplicateAction::NotDuplicate;
};
if row_indices.len() != 1 {
return DuplicateAction::Ambiguous;
}
let row = &self.rows[row_indices[0]];
let needs_repair = (item.session_id.is_some() && !row.has_session_id)
|| (item.hostname.is_some() && !row.has_hostname)
|| (item.cwd.is_some() && !row.has_cwd)
|| (item.duration.is_some() && !row.has_duration)
|| (item.exit_status.is_some() && !row.has_exit_status)
|| (item.more_info.is_some() && matches!(row.metadata, MetadataState::Null));
if needs_repair {
return DuplicateAction::Repair(row.id);
}
if matches!(row.metadata, MetadataState::Malformed) && item.more_info.is_some() {
DuplicateAction::Malformed
} else {
DuplicateAction::Skip
}
}
fn mark_repaired(&mut self, id: HistoryItemId, source: &HistoryItem<HistoryExtraInfo>) {
let Some(row) = self.rows.iter_mut().find(|row| row.id == id) else {
return;
};
row.has_session_id |= source.session_id.is_some();
row.has_hostname |= source.hostname.is_some();
row.has_cwd |= source.cwd.is_some();
row.has_duration |= source.duration.is_some();
row.has_exit_status |= source.exit_status.is_some();
if source.more_info.is_some() && matches!(row.metadata, MetadataState::Null) {
row.metadata = MetadataState::Valid;
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum ImportTarget {
R,
Shell,
}
#[derive(Debug, PartialEq, Eq)]
pub(crate) enum EntryPlan {
Insert(ImportTarget),
Repair {
target: ImportTarget,
id: HistoryItemId,
},
Duplicate,
SkipEmpty,
SkipUnsupported {
mode: String,
},
AmbiguousRepair,
MalformedExistingMetadata,
}
pub(crate) fn plan_entry(
entry: &ImportEntry,
r_dedup: Option<&DedupSet>,
shell_dedup: Option<&DedupSet>,
) -> EntryPlan {
if entry.item.command_line.trim().is_empty() {
return EntryPlan::SkipEmpty;
}
let target = match &entry.mode {
ImportMode::R | ImportMode::Browse | ImportMode::Unspecified => ImportTarget::R,
ImportMode::Shell => ImportTarget::Shell,
ImportMode::Unsupported(mode) => {
return EntryPlan::SkipUnsupported { mode: mode.clone() };
}
};
let dedup = match target {
ImportTarget::R => r_dedup,
ImportTarget::Shell => shell_dedup,
};
let Some(dedup) = dedup else {
return EntryPlan::Insert(target);
};
match dedup.duplicate_action(
&entry.item.command_line,
entry.item.start_timestamp.as_ref(),
&entry.item,
) {
DuplicateAction::NotDuplicate => EntryPlan::Insert(target),
DuplicateAction::Skip => EntryPlan::Duplicate,
DuplicateAction::Repair(id) => EntryPlan::Repair { target, id },
DuplicateAction::Ambiguous => EntryPlan::AmbiguousRepair,
DuplicateAction::Malformed => EntryPlan::MalformedExistingMetadata,
}
}
fn add_plan_warning(result: &mut ImportResult, entry: &ImportEntry, plan: &EntryPlan) {
let command = &entry.item.command_line;
match plan {
EntryPlan::SkipUnsupported { mode } => {
let preview: String = command.chars().take(30).collect();
result
.warnings
.push(format!("Skipped unknown mode '{}': {}...", mode, preview));
}
EntryPlan::AmbiguousRepair => result.warnings.push(format!(
"Could not repair duplicate command '{}': matches multiple rows",
command
)),
EntryPlan::MalformedExistingMetadata => result.warnings.push(format!(
"Could not repair duplicate command '{}': existing metadata is malformed; leaving it unchanged",
command
)),
_ => {}
}
}
fn record_plan(result: &mut ImportResult, entry: &ImportEntry, plan: &EntryPlan) {
match plan {
EntryPlan::Insert(ImportTarget::R) => result.r_imported += 1,
EntryPlan::Insert(ImportTarget::Shell) => result.shell_imported += 1,
EntryPlan::Repair { .. } => result.duplicates_repaired += 1,
EntryPlan::Duplicate => result.duplicates_skipped += 1,
EntryPlan::SkipEmpty => result.skipped += 1,
EntryPlan::SkipUnsupported { .. }
| EntryPlan::AmbiguousRepair
| EntryPlan::MalformedExistingMetadata => {
result.skipped += usize::from(matches!(plan, EntryPlan::SkipUnsupported { .. }));
if !matches!(plan, EntryPlan::SkipUnsupported { .. }) {
result.duplicates_skipped += 1;
}
}
}
add_plan_warning(result, entry, plan);
}
pub fn import_entries_dry_run(
entries: &[ImportEntry],
r_dedup: Option<&DedupSet>,
shell_dedup: Option<&DedupSet>,
) -> ImportResult {
let mut r_dedup = r_dedup.cloned();
let mut shell_dedup = shell_dedup.cloned();
let mut result = ImportResult::default();
for entry in entries {
let plan = plan_entry(entry, r_dedup.as_ref(), shell_dedup.as_ref());
record_plan(&mut result, entry, &plan);
if let EntryPlan::Repair { target, id } = plan {
match target {
ImportTarget::R => r_dedup
.as_mut()
.expect("repair requires an R dedup set")
.mark_repaired(id, &entry.item),
ImportTarget::Shell => shell_dedup
.as_mut()
.expect("repair requires a shell dedup set")
.mark_repaired(id, &entry.item),
}
}
}
result
}
pub fn import_entries(
targets: &mut ImportTargets,
entries: Vec<ImportEntry>,
hostname_override: Option<&str>,
skip_duplicates: bool,
) -> Result<ImportResult> {
let (r_dedup, shell_dedup) = if skip_duplicates {
(
Some(DedupSet::from_history(&targets.r_history)?),
Some(DedupSet::from_history(&targets.shell_history)?),
)
} else {
(None, None)
};
import_entries_with_dedup_sets(targets, entries, hostname_override, r_dedup, shell_dedup)
}
fn import_entries_with_dedup_sets(
targets: &mut ImportTargets,
entries: Vec<ImportEntry>,
hostname_override: Option<&str>,
mut r_dedup: Option<DedupSet>,
mut shell_dedup: Option<DedupSet>,
) -> Result<ImportResult> {
let mut result = ImportResult::default();
for mut entry in entries {
if let Some(hostname) = hostname_override {
entry.item.hostname = Some(hostname.to_owned());
}
let plan = plan_entry(&entry, r_dedup.as_ref(), shell_dedup.as_ref());
match plan {
EntryPlan::Insert(target) => {
let mut item = entry.item;
item.id = None;
let save_result = match target {
ImportTarget::R => targets.r_history.save_imported(item),
ImportTarget::Shell => targets.shell_history.save_imported(item),
};
match save_result {
Ok(_) => match target {
ImportTarget::R => result.r_imported += 1,
ImportTarget::Shell => result.shell_imported += 1,
},
Err(error) => {
result
.warnings
.push(format!("Failed to import entry: {}", error));
result.skipped += 1;
}
}
}
EntryPlan::Repair { target, id } => {
let command = entry.item.command_line.clone();
let store = match target {
ImportTarget::R => &targets.r_history,
ImportTarget::Shell => &targets.shell_history,
};
let source = entry.item;
match store.set_missing_fields_if_empty(id, source.clone()) {
Ok(true) => {
match target {
ImportTarget::R => r_dedup
.as_mut()
.expect("repair requires an R dedup set")
.mark_repaired(id, &source),
ImportTarget::Shell => shell_dedup
.as_mut()
.expect("repair requires a shell dedup set")
.mark_repaired(id, &source),
}
result.duplicates_repaired += 1
}
Ok(false) => result.duplicates_skipped += 1,
Err(error) => {
result.warnings.push(format!(
"Failed to repair duplicate '{}': {}",
command, error
));
result.duplicates_skipped += 1;
}
}
}
other => record_plan(&mut result, &entry, &other),
}
}
Ok(result)
}
pub fn validate_table_name(name: &str) -> Result<()> {
if name.is_empty() {
bail!("Table name cannot be empty");
}
if !name.chars().all(|c| c.is_ascii_alphanumeric() || c == '_') {
bail!(
"Invalid table name '{}': must contain only alphanumeric characters and underscores",
name
);
}
if name.chars().next().is_some_and(|c| c.is_ascii_digit()) {
bail!("Invalid table name '{}': cannot start with a digit", name);
}
if !name.chars().any(|c| c.is_ascii_alphanumeric()) {
bail!(
"Invalid table name '{}': must contain at least one alphanumeric character",
name
);
}
Ok(())
}
pub fn parse_unified_arf_history(
path: &Path,
r_table: &str,
shell_table: &str,
) -> Result<ParsedImport> {
use rusqlite::{Connection, OpenFlags};
validate_table_name(r_table)?;
validate_table_name(shell_table)?;
if r_table == shell_table {
bail!(
"R table name and shell table name must be different (both are '{}')",
r_table
);
}
if !path.exists() {
bail!("arf export file not found: {}", path.display());
}
let db = Connection::open_with_flags(path, OpenFlags::SQLITE_OPEN_READ_ONLY)
.with_context(|| format!("Failed to open arf export file: {}", path.display()))?;
let has_r_table = table_exists(&db, r_table)?;
let has_shell_table = table_exists(&db, shell_table)?;
if !has_r_table && !has_shell_table {
bail!(
"File '{}' does not look like an arf export: missing configured history tables '{}' and '{}'",
path.display(),
r_table,
shell_table
);
}
let mut parsed = ParsedImport::default();
if has_r_table {
let r_entries = read_history_table(&db, path, r_table, ImportMode::R)?;
parsed.entries.extend(r_entries.entries);
parsed.warnings.extend(r_entries.warnings);
}
if has_shell_table {
let shell_entries = read_history_table(&db, path, shell_table, ImportMode::Shell)?;
parsed.entries.extend(shell_entries.entries);
parsed.warnings.extend(shell_entries.warnings);
}
Ok(parsed)
}
fn table_exists(db: &rusqlite::Connection, table_name: &str) -> Result<bool> {
let count: i32 = db
.query_row(
"SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name=?",
[table_name],
|row| row.get(0),
)
.context("Failed to check if table exists")?;
Ok(count > 0)
}
fn read_history_table(
db: &rusqlite::Connection,
source_path: &Path,
table_name: &str,
mode: ImportMode,
) -> Result<ParsedImport> {
use chrono::TimeZone;
let columns = HistoryTableColumns::read(db, table_name)?;
let query = format!(
"SELECT id, command_line, start_timestamp, {}, {}, {}, {}, {}, {} FROM \"{}\" ORDER BY id",
columns.expression("session_id"),
columns.expression("hostname"),
columns.expression("cwd"),
columns.expression("duration_ms"),
columns.expression("exit_status"),
columns.expression("more_info"),
table_name
);
let mut stmt = db.prepare(&query).with_context(|| {
format!(
"Failed to query table '{}' (not a valid history table?)",
table_name
)
})?;
let rows = stmt
.query_map([], |row| {
let id: i64 = row.get(0)?;
let command: String = row.get(1)?;
let ts_millis: Option<i64> = row.get(2)?;
let session_id: Option<i64> = row.get(3)?;
let hostname: Option<String> = row.get(4)?;
let cwd: Option<String> = row.get(5)?;
let duration_millis: Option<i64> = row.get(6)?;
let exit_status: Option<i64> = row.get(7)?;
let raw_metadata: Option<String> = row.get(8)?;
Ok((
id,
command,
ts_millis,
session_id,
hostname,
cwd,
duration_millis,
exit_status,
raw_metadata,
))
})
.context("Failed to query history")?;
let mut parsed = ParsedImport::default();
for row in rows {
let (
id,
command,
ts_millis,
raw_session_id,
hostname,
cwd,
duration_millis,
exit_status,
raw_metadata,
) = row.context("Failed to read history row")?;
let timestamp = ts_millis.and_then(|ms| Utc.timestamp_millis_opt(ms).single());
let session_id = raw_session_id
.map(|id| serde_json::from_str::<HistorySessionId>(&id.to_string()))
.transpose()
.with_context(|| format!("Invalid session_id in row {}", id))?;
let duration = match duration_millis {
Some(ms) if ms >= 0 => Some(Duration::from_millis(ms as u64)),
Some(ms) => {
parsed.warnings.push(format!(
"Invalid negative duration {} for row {} from '{}'; importing with NULL duration",
ms,
id,
source_path.display()
));
None
}
None => None,
};
let metadata = parse_row_metadata(
raw_metadata.as_deref(),
source_path,
HistoryItemId::new(id),
&mut parsed.warnings,
);
parsed.entries.push(ImportEntry {
mode: mode.clone(),
item: HistoryItem {
id: None,
start_timestamp: timestamp,
command_line: command,
session_id,
hostname,
cwd,
duration,
exit_status,
more_info: metadata,
},
});
}
Ok(parsed)
}
struct HistoryTableColumns {
names: HashSet<String>,
}
impl HistoryTableColumns {
fn read(db: &rusqlite::Connection, table_name: &str) -> Result<Self> {
let mut names = HashSet::new();
let query = format!("PRAGMA table_info(\"{}\")", table_name);
let mut stmt = db
.prepare(&query)
.context("Failed to inspect history table")?;
let columns = stmt.query_map([], |row| row.get::<_, String>(1))?;
for column in columns {
names.insert(column?);
}
Ok(Self { names })
}
fn expression(&self, column: &str) -> String {
if column == "more_info" {
if self.names.contains(column) {
"CAST(more_info AS TEXT)".to_string()
} else {
"NULL".to_string()
}
} else if self.names.contains(column) {
column.to_string()
} else {
"NULL".to_string()
}
}
}
fn parse_row_metadata(
raw_metadata: Option<&str>,
source: &Path,
id: HistoryItemId,
warnings: &mut Vec<String>,
) -> Option<HistoryExtraInfo> {
let raw_metadata = raw_metadata?;
match serde_json::from_str::<HistoryExtraInfo>(raw_metadata) {
Ok(metadata) => Some(metadata),
Err(error) => {
warnings.push(format!(
"Could not deserialize metadata for row {} from '{}': {}; importing with NULL metadata",
id.0,
source.display(),
error
));
None
}
}
}
#[cfg(test)]
mod tests;