pub mod app;
pub mod command_palette;
pub mod event;
pub mod qa_protocol;
pub mod qa_ui;
pub mod qa_watcher;
pub mod quit_summary;
pub mod ui;
use std::io::{self, stdout};
use std::pin::Pin;
use std::sync::{Arc, Mutex};
use crossterm::execute;
use crossterm::terminal::{
EnterAlternateScreen, LeaveAlternateScreen, disable_raw_mode, enable_raw_mode,
};
use ratatui::Terminal;
use ratatui::backend::CrosstermBackend;
use tokio::sync::mpsc;
use crate::engine::engine::ExecutionEngine;
use crate::engine::types::{
ExecutionEvent, HarnessDispatchFn, HarnessExecContext, OnEventCallback, OnSaveCallback,
};
use crate::parser::ast::WorkflowDecl;
use crate::state::harness_events::HarnessEventCollector;
use crate::state::types::HarnessEvent;
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 harness_events = Arc::new(Mutex::new(HarnessEventCollector::new(-1, Vec::new())));
let event_tx = tx.clone();
let harness_events_for_engine = harness_events.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));
}
if let Ok(mut collector) = harness_events_for_engine.lock() {
collector.observe_engine_event(event);
}
let _ = event_tx.send(AppMessage::EngineEvent(event.clone()));
});
let on_save: OnSaveCallback = Box::new({
let project_root = project_root_owned.clone();
let harness_events = harness_events.clone();
move |state| {
let harness_event_log = harness_events
.lock()
.ok()
.map(|collector| collector.snapshot());
let persisted = crate::state::adapter::to_persisted_state(state, harness_event_log);
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(),
harness_events.clone(),
tx.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()?;
if app.quit_requested {
let run_id = app
.events
.last()
.map(extract_run_id)
.unwrap_or_else(|| "unknown".to_string());
let status = match app.mode {
app::AppMode::Paused => crate::state::types::RunStatus::Paused,
app::AppMode::Completed => crate::state::types::RunStatus::Completed,
app::AppMode::Failed => crate::state::types::RunStatus::Failed,
_ => crate::state::types::RunStatus::Running,
};
println!(
"{}",
quit_summary::build_quit_summary(
&run_id,
&app.workflow_name,
&status,
&project_root_owned,
)
);
}
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>>>,
harness_events: Arc<Mutex<HarnessEventCollector>>,
tx: mpsc::UnboundedSender<AppMessage>,
) -> 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, context: HarnessExecContext| {
let root = project_root.clone();
let block = exec_block.clone();
let config = config.clone();
let run_id = run_id_holder.clone();
let harness_events = harness_events.clone();
let tx = tx.clone();
Box::pin(async move {
let ctx = crate::harness::RunContext {
run_id,
project_root: root.clone(),
};
let harness_context = crate::harness::HarnessExecutionContext {
run_id: context.run_id,
step_path: context.step_path,
exec_ordinal: context.exec_ordinal,
on_harness_event: Some(Arc::new(move |event: HarnessEvent| {
if let Ok(mut collector) = harness_events.lock() {
let persisted = collector.record(event);
let _ = tx.send(AppMessage::HarnessEvent(persisted));
}
})),
};
crate::harness::dispatch_exec(
&config,
&root,
&block,
Some(&ctx),
Some(&harness_context),
)
.await
})
as Pin<
Box<
dyn std::future::Future<
Output = Result<crate::harness::types::ExecResult, String>,
> + Send,
>,
>
}))
}