pub mod app;
pub mod command_palette;
pub mod event;
pub mod qa_protocol;
pub mod qa_ui;
pub mod qa_watcher;
pub mod ui;
use std::io::{self, stdout};
use std::pin::Pin;
use std::sync::Arc;
use crossterm::execute;
use crossterm::terminal::{
disable_raw_mode, enable_raw_mode, EnterAlternateScreen, LeaveAlternateScreen,
};
use ratatui::backend::CrosstermBackend;
use ratatui::Terminal;
use tokio::sync::mpsc;
use crate::engine::engine::ExecutionEngine;
use crate::engine::types::{ExecutionEvent, HarnessDispatchFn, OnEventCallback, OnSaveCallback};
use crate::parser::ast::WorkflowDecl;
use self::app::{App, AppMessage};
pub async fn run_tui(
workflows: Vec<WorkflowDecl>,
project_root: &str,
workflow_name: &str,
) -> Result<(), Box<dyn std::error::Error>> {
enable_raw_mode()?;
let mut stdout = stdout();
execute!(stdout, EnterAlternateScreen)?;
let backend = CrosstermBackend::new(stdout);
let mut terminal = Terminal::new(backend)?;
let panic_hook = std::panic::take_hook();
std::panic::set_hook(Box::new(move |info| {
let _ = disable_raw_mode();
let _ = execute!(io::stdout(), LeaveAlternateScreen);
panic_hook(info);
}));
let project_root_owned = project_root.to_string();
let (tx, rx) = mpsc::unbounded_channel::<AppMessage>();
let run_id_holder: Arc<std::sync::RwLock<Option<String>>> =
Arc::new(std::sync::RwLock::new(None));
let run_id_for_event = run_id_holder.clone();
let run_id_for_submit = run_id_holder.clone();
let event_tx = tx.clone();
let on_event: OnEventCallback = Box::new(move |event: &ExecutionEvent| {
if let Ok(mut guard) = run_id_for_event.write()
&& guard.is_none()
{
*guard = Some(extract_run_id(event));
}
let _ = event_tx.send(AppMessage::EngineEvent(event.clone()));
});
let on_save: OnSaveCallback = Box::new({
let project_root = project_root_owned.clone();
move |state| {
let persisted = crate::state::adapter::to_persisted_state(state);
if let Err(e) = crate::state::persistence::save_state(&persisted, &project_root) {
eprintln!("[o7] Warning: failed to save state: {}", e);
}
}
});
let dispatch = make_tui_dispatch_fn(project_root_owned.clone(), run_id_holder.clone())
.map_err(|e| -> Box<dyn std::error::Error> { e.into() })?;
let mut engine = ExecutionEngine::new(workflows.clone(), dispatch, Some(on_event), Some(on_save));
engine.suspend_on_pause = true;
let pause_flag = engine.pause_requested.clone();
let pause_notify = engine.pause_notify.clone();
let mut app = App::new(&workflows, workflow_name, pause_flag.clone(), pause_notify);
let wf_name = workflow_name.to_string();
let engine_tx = tx.clone();
std::thread::spawn(move || {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("Failed to create engine runtime");
rt.block_on(async move {
let (_run_id, _events) = engine.start_run(&wf_name).await;
let _ = engine_tx.send(AppMessage::EngineFinished);
});
});
let qa_watcher = None::<self::qa_watcher::QAWatcher>;
let qa_project_root = project_root_owned.clone();
let on_qa_submit: event::QASubmitCallback = Box::new(move |seq, answers_file| {
if let Ok(guard) = run_id_for_submit.read()
&& let Some(ref run_id) = *guard
{
let qa_dir = format!("{}/.7/runs/{}/qa", qa_project_root, run_id);
if let Err(e) = std::fs::create_dir_all(&qa_dir) {
eprintln!("[o7] Warning: failed to create qa dir: {}", e);
return;
}
let filename = format!(
"q-answers-{}.json",
qa_protocol::format_seq(seq)
);
let filepath = format!("{}/{}", qa_dir, filename);
let tmp_filepath = format!("{}.tmp", filepath);
match serde_json::to_string_pretty(answers_file) {
Ok(json) => {
if let Err(e) = std::fs::write(&tmp_filepath, &json) {
eprintln!("[o7] Warning: failed to write answer file: {}", e);
return;
}
if let Err(e) = std::fs::rename(&tmp_filepath, &filepath) {
eprintln!("[o7] Warning: failed to rename answer file: {}", e);
}
}
Err(e) => {
eprintln!("[o7] Warning: failed to serialize answers: {}", e);
}
}
}
});
let result = event::run_event_loop(
&mut terminal,
&mut app,
rx,
qa_watcher,
Some(on_qa_submit),
run_id_holder,
project_root_owned.clone(),
)
.await;
disable_raw_mode()?;
execute!(terminal.backend_mut(), LeaveAlternateScreen)?;
terminal.show_cursor()?;
result.map_err(|e| e.into())
}
fn extract_run_id(event: &ExecutionEvent) -> String {
match event {
ExecutionEvent::StepStarted { runId, .. }
| ExecutionEvent::StepCompleted { runId, .. }
| ExecutionEvent::StepFailed { runId, .. }
| ExecutionEvent::BranchStarted { runId, .. }
| ExecutionEvent::BranchCompleted { runId, .. }
| ExecutionEvent::BranchFailed { runId, .. }
| ExecutionEvent::JoinStarted { runId, .. }
| ExecutionEvent::RunCompleted { runId }
| ExecutionEvent::RunFailed { runId, .. }
| ExecutionEvent::RunPaused { runId, .. }
| ExecutionEvent::CheckEvaluated { runId, .. }
| ExecutionEvent::MatchEvaluated { runId, .. }
| ExecutionEvent::SafeBoundary { runId, .. } => runId.clone(),
}
}
fn make_tui_dispatch_fn(
project_root: String,
run_id_holder: Arc<std::sync::RwLock<Option<String>>>,
) -> Result<HarnessDispatchFn, String> {
let config = match crate::harness::config::load_harness_config(&project_root) {
Ok(c) => c,
Err(_) => crate::harness::types::HarnessConfig { harness: std::collections::HashMap::new() },
};
let config = Arc::new(config);
Ok(Arc::new(move |exec_block| {
let root = project_root.clone();
let block = exec_block.clone();
let config = config.clone();
let run_id = run_id_holder.clone();
Box::pin(async move {
let ctx = crate::harness::RunContext {
run_id,
project_root: root.clone(),
};
crate::harness::dispatch_exec(&config, &root, &block, Some(&ctx)).await
})
as Pin<
Box<
dyn std::future::Future<
Output = Result<crate::harness::types::ExecResult, String>,
> + Send,
>,
>
}))
}