use std::collections::HashSet;
use std::path::Path;
#[derive(Debug, Clone)]
pub struct RefinePlan {
pub completed: HashSet<(String, String)>,
pub seen_identities: HashSet<(String, String)>,
pub completed_hashes: std::collections::HashMap<(String, String), PriorCompletion>,
pub next_exec_id: u64,
pub prior_outcomes_seen: usize,
pub scope: RefineScope,
}
#[derive(Debug, Clone, Default)]
pub struct PriorCompletion {
pub phase_hash: Option<String>,
pub params_consumed: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum SkipBlocker {
NoPrior,
BaseChanged,
ParamChanged(String),
}
impl std::fmt::Display for SkipBlocker {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::NoPrior => write!(f, "no comparable prior outcome"),
Self::BaseChanged => write!(f, "scope or phase config changed"),
Self::ParamChanged(name) => write!(f, "param '{name}' changed"),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RefineScope {
Missing,
Changed,
All,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ExecutionQualifier {
Specific(u64),
All,
}
pub const LATEST_LITERAL: &str = "latest";
pub fn warn_multi_execution_default(session_dir: &std::path::Path) {
let db_path = session_dir.join("metrics.db");
if !db_path.exists() {
return;
}
let conn = match rusqlite::Connection::open_with_flags(
&db_path,
rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
) {
Ok(c) => c,
Err(_) => return,
};
let exists: bool = conn
.query_row(
"SELECT EXISTS(SELECT 1 FROM sqlite_master \
WHERE type='table' AND name='executions')",
[],
|r| r.get::<_, i64>(0),
)
.map(|n| n != 0)
.unwrap_or(false);
if !exists {
return;
}
let mut stmt = match conn.prepare(
"SELECT exec_id, verb, scope, disposition, started_at_nanos \
FROM executions ORDER BY exec_id DESC LIMIT 3",
) {
Ok(s) => s,
Err(_) => return,
};
type ExecRow = (i64, String, Option<String>, Option<String>, i64);
let rows: Vec<ExecRow> = match stmt.query_map([], |r| {
Ok((
r.get::<_, i64>(0)?,
r.get::<_, String>(1)?,
r.get::<_, Option<String>>(2)?,
r.get::<_, Option<String>>(3)?,
r.get::<_, i64>(4)?,
))
}) {
Ok(it) => it.filter_map(Result::ok).collect(),
Err(_) => return,
};
let total: i64 = conn
.query_row("SELECT COUNT(*) FROM executions", [], |r| r.get(0))
.unwrap_or(0);
if total < 2 || rows.is_empty() {
return;
}
eprintln!(
"session has {total} execution(s); implicit qualifier `exec_id=latest` \
resolved to exec_id={latest}. recent (newest first):",
latest = rows[0].0,
);
for (i, (exec_id, verb, scope, disposition, _)) in rows.iter().enumerate() {
let marker = if i == 0 { " <-- latest" } else { "" };
let scope_part = scope
.as_deref()
.map(|s| format!(" scope={s}"))
.unwrap_or_default();
let disp_part = disposition
.as_deref()
.map(|d| format!(" {d}"))
.unwrap_or_else(|| " (in-flight)".to_string());
eprintln!(" exec_id={exec_id} verb={verb}{scope_part}{disp_part}{marker}");
}
eprintln!(" (pass `--execution=<n>` to target one, `--all-executions` to aggregate)");
}
impl ExecutionQualifier {
pub fn specific(n: u64) -> Self {
Self::Specific(n)
}
pub fn all() -> Self {
Self::All
}
pub fn latest(session_dir: &std::path::Path) -> Self {
match latest_exec_id_for_session(session_dir) {
Some(n) => Self::Specific(n),
None => Self::Specific(1),
}
}
pub fn matches_all(&self) -> bool {
matches!(self, Self::All)
}
pub fn specific_id(&self) -> Option<u64> {
match self {
Self::Specific(n) => Some(*n),
Self::All => None,
}
}
}
pub fn latest_exec_id_for_session(session_dir: &std::path::Path) -> Option<u64> {
let db_path = session_dir.join("metrics.db");
if !db_path.exists() {
return None;
}
let conn = rusqlite::Connection::open_with_flags(
&db_path,
rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
)
.ok()?;
let table_exists: bool = conn
.query_row(
"SELECT EXISTS(SELECT 1 FROM sqlite_master \
WHERE type='table' AND name='phase_outcomes')",
[],
|r| r.get::<_, i64>(0),
)
.map(|n| n != 0)
.unwrap_or(false);
if !table_exists {
return None;
}
let max: Option<i64> = conn
.query_row("SELECT MAX(exec_id) FROM phase_outcomes", [], |r| {
r.get::<_, Option<i64>>(0)
})
.ok()
.flatten();
max.map(|v| v.max(0) as u64)
}
impl RefinePlan {
pub fn is_completed(&self, phase_name: &str, phase_labels: &str) -> bool {
self.completed
.contains(&(phase_name.to_string(), phase_labels.to_string()))
}
pub fn is_unchanged(
&self,
phase_name: &str,
phase_labels: &str,
current_hex: &str,
current_params: &std::collections::HashMap<String, String>,
) -> bool {
self.unchanged_verdict(phase_name, phase_labels, current_hex, current_params)
.is_ok()
}
pub fn unchanged_verdict(
&self,
phase_name: &str,
phase_labels: &str,
current_hex: &str,
current_params: &std::collections::HashMap<String, String>,
) -> Result<(), SkipBlocker> {
let key = (phase_name.to_string(), phase_labels.to_string());
let Some(prior) = self.completed_hashes.get(&key) else {
return Err(SkipBlocker::NoPrior);
};
match prior.phase_hash.as_deref() {
None => return Err(SkipBlocker::NoPrior),
Some(prior_hex) if prior_hex != current_hex => return Err(SkipBlocker::BaseChanged),
Some(_) => {}
}
let Some(json) = prior.params_consumed.as_deref() else {
return Err(SkipBlocker::NoPrior);
};
let Ok(stored) = serde_json::from_str::<std::collections::BTreeMap<String, String>>(json)
else {
return Err(SkipBlocker::NoPrior);
};
for (name, stored_digest) in stored {
let current = current_params
.get(&name)
.map(|v| crate::checkpoint::params_scope::value_digest(v));
if current.as_deref() != Some(stored_digest.as_str()) {
return Err(SkipBlocker::ParamChanged(name));
}
}
Ok(())
}
pub fn should_skip(
&self,
phase_name: &str,
phase_labels: &str,
current_hex: &str,
current_params: &std::collections::HashMap<String, String>,
) -> bool {
match self.scope {
RefineScope::All => false,
RefineScope::Missing => self.is_completed(phase_name, phase_labels),
RefineScope::Changed => {
self.is_unchanged(phase_name, phase_labels, current_hex, current_params)
}
}
}
pub fn load_from_session_dir(session_dir: &Path) -> Option<Self> {
let db_path = session_dir.join("metrics.db");
if !db_path.exists() {
return None;
}
let conn = match rusqlite::Connection::open_with_flags(
&db_path,
rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
) {
Ok(c) => c,
Err(e) => {
crate::diag!(
crate::observer::LogLevel::Warn,
"refine: failed to open {}: {e}",
db_path.display()
);
return None;
}
};
let exists: bool = conn
.query_row(
"SELECT EXISTS(SELECT 1 FROM sqlite_master \
WHERE type='table' AND name='phase_outcomes')",
[],
|r| r.get::<_, i64>(0),
)
.map(|n| n != 0)
.unwrap_or(false);
if !exists {
return None;
}
let has_params_col: bool = conn
.prepare("PRAGMA table_info(phase_outcomes)")
.ok()
.and_then(|mut s| {
let mut found = false;
let mut rows = s.query([]).ok()?;
while let Ok(Some(r)) = rows.next() {
if r.get::<_, String>(1)
.map(|n| n == "params_consumed")
.unwrap_or(false)
{
found = true;
}
}
Some(found)
})
.unwrap_or(false);
let pc_col = if has_params_col {
"params_consumed"
} else {
"NULL"
};
let mut stmt = match conn.prepare(&format!(
"SELECT exec_id, phase_name, phase_labels, status, phase_hash, \
{pc_col}, ended_at_nanos \
FROM phase_outcomes \
ORDER BY ended_at_nanos"
)) {
Ok(s) => s,
Err(e) => {
crate::diag!(
crate::observer::LogLevel::Warn,
"refine: failed to prepare query: {e}"
);
return None;
}
};
let rows = match stmt.query_map([], |row| {
Ok((
row.get::<_, i64>(0)?,
row.get::<_, String>(1)?,
row.get::<_, String>(2)?,
row.get::<_, String>(3)?,
row.get::<_, Option<String>>(4)?,
row.get::<_, Option<String>>(5)?,
))
}) {
Ok(r) => r,
Err(e) => {
crate::diag!(
crate::observer::LogLevel::Warn,
"refine: failed to query phase_outcomes: {e}"
);
return None;
}
};
let mut completed = HashSet::new();
let mut seen_identities: HashSet<(String, String)> = HashSet::new();
let mut completed_hashes: std::collections::HashMap<(String, String), PriorCompletion> =
std::collections::HashMap::new();
let mut max_exec_id: u64 = 0;
let mut count: usize = 0;
for row in rows.flatten() {
let (exec_id, name, labels, status, phase_hash, params_consumed) = row;
let exec_id = exec_id.max(0) as u64;
if exec_id > max_exec_id {
max_exec_id = exec_id;
}
count += 1;
seen_identities.insert((name.clone(), labels.clone()));
if status == "completed" {
completed.insert((name.clone(), labels.clone()));
completed_hashes.insert(
(name, labels),
PriorCompletion {
phase_hash,
params_consumed,
},
);
}
}
let executions_max: u64 = conn
.query_row(
"SELECT EXISTS(SELECT 1 FROM sqlite_master \
WHERE type='table' AND name='executions')",
[],
|r| r.get::<_, i64>(0),
)
.map(|n| n != 0)
.ok()
.filter(|exists| *exists)
.and_then(|_| {
conn.query_row("SELECT MAX(exec_id) FROM executions", [], |r| {
r.get::<_, Option<i64>>(0)
})
.ok()
.flatten()
})
.map(|v| v.max(0) as u64)
.unwrap_or(0);
let max_exec_id = max_exec_id.max(executions_max);
Some(Self {
completed,
seen_identities,
completed_hashes,
next_exec_id: max_exec_id + 1,
prior_outcomes_seen: count,
scope: RefineScope::Missing,
})
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn missing_db_returns_none() {
let tmp = tempfile::tempdir().unwrap();
assert!(RefinePlan::load_from_session_dir(tmp.path()).is_none());
}
#[test]
fn db_without_phase_outcomes_returns_none() {
let tmp = tempfile::tempdir().unwrap();
let db_path = tmp.path().join("metrics.db");
let conn = rusqlite::Connection::open(&db_path).unwrap();
conn.execute("CREATE TABLE other (x INTEGER)", []).unwrap();
drop(conn);
assert!(RefinePlan::load_from_session_dir(tmp.path()).is_none());
}
fn make_db_with_outcomes(dir: &Path, rows: &[(u64, &str, &str, &str)]) {
let db_path = dir.join("metrics.db");
let conn = rusqlite::Connection::open(&db_path).unwrap();
conn.execute(
"CREATE TABLE phase_outcomes (
session TEXT NOT NULL,
exec_id INTEGER NOT NULL,
phase_name TEXT NOT NULL,
phase_labels TEXT NOT NULL,
status TEXT NOT NULL,
duration_secs REAL NOT NULL DEFAULT 0,
started_at_nanos INTEGER NOT NULL DEFAULT 0,
ended_at_nanos INTEGER NOT NULL DEFAULT 0,
phase_hash TEXT,
PRIMARY KEY (session, exec_id, phase_name, phase_labels)
)",
[],
)
.unwrap();
for (exec, name, labels, status) in rows {
conn.execute(
"INSERT INTO phase_outcomes (session, exec_id, phase_name, phase_labels, status) \
VALUES ('s', ?1, ?2, ?3, ?4)",
rusqlite::params![*exec as i64, name, labels, status],
)
.unwrap();
}
}
#[test]
fn completed_phases_populate_skip_set() {
let tmp = tempfile::tempdir().unwrap();
make_db_with_outcomes(
tmp.path(),
&[
(1, "schema", "", "completed"),
(1, "load_data", "k=10", "completed"),
(1, "query", "k=10,limit=20", "failed"),
],
);
let plan = RefinePlan::load_from_session_dir(tmp.path()).unwrap();
assert_eq!(plan.next_exec_id, 2);
assert_eq!(plan.prior_outcomes_seen, 3);
assert!(plan.is_completed("schema", ""));
assert!(plan.is_completed("load_data", "k=10"));
assert!(!plan.is_completed("query", "k=10,limit=20"));
}
#[test]
fn next_exec_id_bumps_past_max_prior() {
let tmp = tempfile::tempdir().unwrap();
make_db_with_outcomes(
tmp.path(),
&[
(1, "schema", "", "completed"),
(3, "query", "", "completed"),
(2, "load", "", "completed"),
],
);
let plan = RefinePlan::load_from_session_dir(tmp.path()).unwrap();
assert_eq!(plan.next_exec_id, 4);
}
#[test]
fn empty_table_yields_exec_id_1() {
let tmp = tempfile::tempdir().unwrap();
make_db_with_outcomes(tmp.path(), &[]);
let plan = RefinePlan::load_from_session_dir(tmp.path()).unwrap();
assert_eq!(plan.next_exec_id, 1);
assert_eq!(plan.prior_outcomes_seen, 0);
assert!(plan.completed.is_empty());
}
#[test]
fn execution_qualifier_specific_carries_id() {
let q = ExecutionQualifier::specific(7);
assert_eq!(q.specific_id(), Some(7));
assert!(!q.matches_all());
}
#[test]
fn execution_qualifier_all_carries_no_id() {
let q = ExecutionQualifier::all();
assert_eq!(q.specific_id(), None);
assert!(q.matches_all());
}
#[test]
fn execution_qualifier_latest_resolves_to_max_exec_id() {
let tmp = tempfile::tempdir().unwrap();
make_db_with_outcomes(
tmp.path(),
&[
(1, "p1", "", "completed"),
(3, "p2", "", "completed"),
(2, "p3", "", "completed"),
],
);
let q = ExecutionQualifier::latest(tmp.path());
assert_eq!(
q.specific_id(),
Some(3),
"latest MUST resolve to max(exec_id)=3"
);
assert!(
!q.matches_all(),
"latest MUST narrow to a specific id, not aggregate"
);
}
#[test]
fn execution_qualifier_latest_on_empty_db_yields_specific_1() {
let tmp = tempfile::tempdir().unwrap();
let q = ExecutionQualifier::latest(tmp.path());
assert_eq!(
q.specific_id(),
Some(1),
"latest on empty db MUST yield Specific(1), not All"
);
}
#[test]
fn execution_qualifier_latest_on_missing_table_yields_specific_1() {
let tmp = tempfile::tempdir().unwrap();
let db_path = tmp.path().join("metrics.db");
let conn = rusqlite::Connection::open(&db_path).unwrap();
conn.execute("CREATE TABLE other (x INTEGER)", []).unwrap();
drop(conn);
let q = ExecutionQualifier::latest(tmp.path());
assert_eq!(q.specific_id(), Some(1));
}
}