hypersteeldb 0.5.2

A database that compiles questions instead of guessing answers: typed vocabulary discovered from your documents, queries type-checked before they run, roaring-bitmap set algebra over reified hyperedges, and Dempster-Shafer evidence with an explicit conflict guard.
Documentation
//! SteelDB TUI — realtime "AI grep" over a folder. A worker thread keeps the text-projection models
//! (SPO tagger + SPLADE) and the local Qwen model HOT and owns a growable corpus: documents are
//! projected file-by-file and the corpus fills in live, then you ask questions or `/add <path>` more —
//! nothing reloads. Retrieval-as-reasoning, grounded in the corpus.
//!
//!   cargo run --features tui,paddock,onnx,docs,ocr --bin tui -- <dir>
//!
//! Local model by default (Paddock → ollama on :11434). STEELDB_LLM=bedrock for Claude.

use ratatui::crossterm::event::{self, Event, KeyCode, KeyModifiers};
use ratatui::crossterm::execute;
use ratatui::crossterm::terminal::{disable_raw_mode, enable_raw_mode, EnterAlternateScreen, LeaveAlternateScreen};
use ratatui::prelude::*;
use ratatui::widgets::{Block, Borders, Paragraph, Wrap};
use std::io::stdout;
use std::path::{Path, PathBuf};
use std::sync::mpsc;
use std::time::Duration;
use steeldb::agent::{run_agent, ProviderConfig};
use steeldb::projectors::{CsvProjector, JsonProjector, JsonlProjector, TextEngine};
use steeldb::{Corpus, CorpusKind, Projector};

const SPINNER: [&str; 8] = ["⠋", "⠙", "⠹", "⠸", "⠼", "⠴", "⠦", "⠧"];

fn provider_config() -> ProviderConfig {
    if std::env::var("STEELDB_LLM").as_deref() == Ok("bedrock") {
        return ProviderConfig::from_env();
    }
    ProviderConfig::Paddock {
        base_url: std::env::var("STEELDB_PADDOCK_URL").unwrap_or_else(|_| "http://localhost:11434/v1".to_string()),
        model: std::env::var("STEELDB_PADDOCK_MODEL").unwrap_or_else(|_| "qwen3:1.7b".to_string()),
        api_key: None,
    }
}

/// Load the hot text-projection engine from env/defaults (models loaded once, reused for every doc).
fn hot_engine() -> Option<TextEngine> {
    let ml = steeldb::paths::model_dir("step0_bundle_ml", "STEELDB_ML_BUNDLE", "spo.onnx")?;
    let splade = steeldb::paths::model_dir("splade", "STEELDB_SPLADE_DIR", "splade.onnx");
    TextEngine::load(&ml, splade.as_deref()).ok()
}

/// Recursively collect ingestable files under a path (or just the file itself).
fn collect(path: &Path) -> Vec<PathBuf> {
    if path.is_file() {
        return vec![path.to_path_buf()];
    }
    let mut out = Vec::new();
    let mut stack = vec![path.to_path_buf()];
    while let Some(d) = stack.pop() {
        let Ok(rd) = std::fs::read_dir(&d) else { continue };
        for e in rd.flatten() {
            let name = e.file_name().to_string_lossy().to_string();
            if name.starts_with('.') || name == "node_modules" || name == "target" {
                continue;
            }
            let p = e.path();
            if p.is_dir() {
                stack.push(p);
            } else {
                out.push(p);
            }
        }
    }
    out.sort();
    out
}

/// Project one file into the live corpus using the hot engine. Returns situations added.
fn ingest_file(corpus: &mut Corpus, engine: &mut Option<TextEngine>, path: &Path) -> usize {
    let ext = path.extension().and_then(|e| e.to_str()).unwrap_or("").to_lowercase();
    let name = path.file_name().map(|n| n.to_string_lossy().to_string()).unwrap_or_default();
    let src_tok = format!("src/{}", steeldb::projector::slug(&name));
    let mut n = 0usize;

    let mut sink = |s: steeldb::Situation, disp: String| {
        let mut toks = s.tokens;
        toks.push(src_tok.clone());
        corpus.add_situation_num(toks, vec![name.clone(), disp], s.numbers);
        n += 1;
    };

    match ext.as_str() {
        "csv" | "tsv" => {
            if let Ok(p) = CsvProjector::open(path) {
                let _ = Box::new(p).project(&mut |s| {
                    let disp = s.display.join(" · ");
                    sink(s, disp);
                });
            }
        }
        "json" | "ndjson" => {
            let _ = Box::new(JsonProjector::open(path)).project(&mut |s| {
                let disp = s.display.join(" · ");
                sink(s, disp);
            });
        }
        "jsonl" => {
            let _ = Box::new(JsonlProjector::open(path, None)).project(&mut |s| {
                let disp = s.display.join(" · ");
                sink(s, disp);
            });
        }
        _ => {
            #[cfg(feature = "docs")]
            if steeldb::docs::is_doc_ext(&ext) {
                if let Some(eng) = engine.as_mut() {
                    if let Ok(Some(text)) = steeldb::docs::extract_text(path) {
                        eng.project_text(&text, &mut |s| {
                            let disp = s.display.join(" ");
                            sink(s, disp);
                        });
                    }
                }
            }
            #[cfg(not(feature = "docs"))]
            let _ = engine;
        }
    }
    n
}

enum Job {
    Ask(String),
    Add(String),
}
enum Reply {
    Status(String),                    // transient status (ingesting <file>…)
    Line(Line<'static>),               // append to transcript
    Stats { sits: u32, toks: usize },  // update header counts
    Done,                              // job finished → clear busy
}

fn main() -> Result<(), Box<dyn std::error::Error>> {
    let args: Vec<String> = std::env::args().collect();
    let dir = args.get(1).cloned();

    let (job_tx, job_rx) = mpsc::channel::<Job>();
    let (rep_tx, rep_rx) = mpsc::channel::<Reply>();

    // worker owns the hot models + provider + the growable corpus; handles one job at a time.
    let cfg = provider_config();
    std::thread::spawn(move || {
        let rt = match tokio::runtime::Builder::new_current_thread().enable_all().build() {
            Ok(rt) => rt,
            Err(e) => {
                let _ = rep_tx.send(Reply::Line(err_line(&format!("runtime: {e}"))));
                return;
            }
        };
        let provider = rt.block_on(cfg.build());
        let mut engine = hot_engine();
        if engine.is_none() {
            let _ = rep_tx.send(Reply::Line(dim_line("text models unavailable — set STEELDB_ML_BUNDLE (docs/text won't project)")));
        }
        let mut corpus = Corpus::new_incremental("(live)", vec!["file".into(), "record".into()], CorpusKind::Csv);

        while let Ok(job) = job_rx.recv() {
            match job {
                Job::Add(path) => {
                    let files = collect(Path::new(&path));
                    if files.is_empty() {
                        let _ = rep_tx.send(Reply::Line(dim_line(&format!("nothing to ingest at {path}"))));
                    }
                    for f in files {
                        let fname = f.file_name().map(|n| n.to_string_lossy().to_string()).unwrap_or_default();
                        let _ = rep_tx.send(Reply::Status(format!("ingesting {fname}…")));
                        let added = ingest_file(&mut corpus, &mut engine, &f);
                        let s = corpus.stats();
                        let _ = rep_tx.send(Reply::Line(dim_line(&format!("+ {fname}  ({added} situations)"))));
                        let _ = rep_tx.send(Reply::Stats { sits: s.situations, toks: s.vocab });
                    }
                }
                Job::Ask(q) => match &provider {
                    Ok(p) => match rt.block_on(run_agent(p.as_ref(), &corpus, &q, 14)) {
                        Ok(ans) => {
                            for l in answer_lines(&ans.answer) {
                                let _ = rep_tx.send(Reply::Line(l));
                            }
                            let tools: Vec<String> = ans.trace.iter().map(|t| t.name.clone()).collect();
                            if !tools.is_empty() {
                                let _ = rep_tx.send(Reply::Line(dim_line(&format!("· {} tool calls: {}", ans.trace.len(), tools.join(", ")))));
                            }
                        }
                        Err(e) => {
                            let _ = rep_tx.send(Reply::Line(err_line(&e)));
                        }
                    },
                    Err(e) => {
                        let _ = rep_tx.send(Reply::Line(err_line(&format!("provider: {e}"))));
                    }
                },
            }
            let _ = rep_tx.send(Reply::Line(Line::from("")));
            let _ = rep_tx.send(Reply::Done);
        }
    });

    // kick off the initial ingest so the corpus fills in live
    let model = match provider_config() {
        ProviderConfig::Paddock { model, .. } => model,
        ProviderConfig::Bedrock { model_id, .. } => model_id,
    };
    let mut app = App::new(model);
    if let Some(d) = &dir {
        app.busy = true;
        app.status = "ingesting…".into();
        let _ = job_tx.send(Job::Add(d.clone()));
        app.push(dim_line(&format!("ingesting {d} …  ask a question once it fills in, or /add <path> for more")));
    } else {
        app.push(dim_line("empty corpus — /add <path> to ingest a file or folder, then ask a question"));
    }
    app.push(dim_line("Enter=send · /add <path> to ingest · PgUp/PgDn scroll · Ctrl-C quit"));

    enable_raw_mode()?;
    let mut out = stdout();
    execute!(out, EnterAlternateScreen)?;
    let mut terminal = Terminal::new(CrosstermBackend::new(out))?;
    let res = run_ui(&mut terminal, &mut app, &job_tx, &rep_rx);
    disable_raw_mode()?;
    execute!(terminal.backend_mut(), LeaveAlternateScreen)?;
    terminal.show_cursor()?;
    res.map_err(Into::into)
}

fn dim_line(s: &str) -> Line<'static> {
    Line::from(Span::styled(s.to_string(), Style::default().fg(Color::DarkGray)))
}
fn err_line(s: &str) -> Line<'static> {
    Line::from(vec![Span::styled("! ".to_string(), Style::default().fg(Color::Red)), Span::raw(s.to_string())])
}
fn answer_lines(text: &str) -> Vec<Line<'static>> {
    text.split('\n')
        .enumerate()
        .map(|(i, seg)| {
            let prefix = if i == 0 { "‹ " } else { "  " };
            Line::from(vec![Span::styled(prefix.to_string(), Style::default().fg(Color::Green)), Span::raw(seg.to_string())])
        })
        .collect()
}

struct App {
    lines: Vec<Line<'static>>,
    input: String,
    busy: bool,
    tick: usize,
    scroll: u16,
    stick_bottom: bool,
    status: String,
    sits: u32,
    toks: usize,
    model: String,
}

impl App {
    fn new(model: String) -> App {
        App { lines: Vec::new(), input: String::new(), busy: false, tick: 0, scroll: 0, stick_bottom: true, status: String::new(), sits: 0, toks: 0, model }
    }
    fn push(&mut self, line: Line<'static>) {
        self.lines.push(line);
        self.stick_bottom = true;
    }
}

fn run_ui<B: Backend>(terminal: &mut Terminal<B>, app: &mut App, job_tx: &mpsc::Sender<Job>, rep_rx: &mpsc::Receiver<Reply>) -> std::io::Result<()> {
    loop {
        while let Ok(rep) = rep_rx.try_recv() {
            match rep {
                Reply::Status(s) => app.status = s,
                Reply::Line(l) => app.push(l),
                Reply::Stats { sits, toks } => {
                    app.sits = sits;
                    app.toks = toks;
                }
                Reply::Done => {
                    app.busy = false;
                    app.status.clear();
                }
            }
        }

        terminal.draw(|f| draw(f, app))?;

        if event::poll(Duration::from_millis(120))? {
            if let Event::Key(k) = event::read()? {
                if k.modifiers.contains(KeyModifiers::CONTROL) && matches!(k.code, KeyCode::Char('c')) {
                    return Ok(());
                }
                match k.code {
                    KeyCode::Esc => return Ok(()),
                    KeyCode::Enter => {
                        let line = app.input.trim().to_string();
                        app.input.clear();
                        if line.is_empty() || app.busy {
                            // ignore
                        } else if let Some(rest) = line.strip_prefix("/add ").or_else(|| line.strip_prefix(":add ")) {
                            app.push(Line::from(Span::styled(format!("/add {}", rest.trim()), Style::default().fg(Color::Yellow))));
                            app.busy = true;
                            app.status = "ingesting…".into();
                            let _ = job_tx.send(Job::Add(rest.trim().to_string()));
                        } else {
                            app.push(Line::from(vec![Span::styled("› ".to_string(), Style::default().fg(Color::Cyan).add_modifier(Modifier::BOLD)), Span::raw(line.clone())]));
                            app.busy = true;
                            app.status = "thinking…".into();
                            let _ = job_tx.send(Job::Ask(line));
                        }
                    }
                    KeyCode::Char(c) => app.input.push(c),
                    KeyCode::Backspace => {
                        app.input.pop();
                    }
                    KeyCode::PageUp => {
                        app.stick_bottom = false;
                        app.scroll = app.scroll.saturating_sub(8);
                    }
                    KeyCode::Up => {
                        app.stick_bottom = false;
                        app.scroll = app.scroll.saturating_sub(1);
                    }
                    KeyCode::PageDown => app.scroll = app.scroll.saturating_add(8),
                    KeyCode::Down => app.scroll = app.scroll.saturating_add(1),
                    _ => {}
                }
            }
        }
        app.tick = app.tick.wrapping_add(1);
    }
}

fn draw(f: &mut Frame, app: &mut App) {
    let chunks = Layout::default()
        .direction(Direction::Vertical)
        .constraints([Constraint::Length(1), Constraint::Min(1), Constraint::Length(3)])
        .split(f.area());

    let header = format!(" SteelDB · {} situations, {} tokens · {} ● hot", app.sits, app.toks, app.model);
    f.render_widget(
        Paragraph::new(Line::from(Span::styled(header, Style::default().fg(Color::Black).bg(Color::Cyan).add_modifier(Modifier::BOLD)))).style(Style::default().bg(Color::Cyan)),
        chunks[0],
    );

    let view_h = chunks[1].height.saturating_sub(2);
    let follow_max = (app.lines.len() as u16).saturating_sub(view_h);
    let scroll = if app.stick_bottom { follow_max } else { app.scroll.min(follow_max) };
    app.scroll = scroll;
    if scroll >= follow_max {
        app.stick_bottom = true;
    }
    f.render_widget(
        Paragraph::new(app.lines.clone()).block(Block::default().borders(Borders::ALL).title(" conversation ")).wrap(Wrap { trim: false }).scroll((scroll, 0)),
        chunks[1],
    );

    let prompt = if app.busy {
        Line::from(vec![
            Span::styled(format!(" {} ", SPINNER[app.tick % SPINNER.len()]), Style::default().fg(Color::Yellow)),
            Span::styled(app.status.clone(), Style::default().fg(Color::DarkGray)),
        ])
    } else {
        Line::from(vec![Span::styled(" › ".to_string(), Style::default().fg(Color::Cyan)), Span::raw(app.input.clone()), Span::styled("▏".to_string(), Style::default().fg(Color::Cyan))])
    };
    f.render_widget(Paragraph::new(prompt).block(Block::default().borders(Borders::ALL).title(" ask ")), chunks[2]);
}