Skip to main content

nmbrs_runtime/
refine_plan.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! SRD-77 refine: the per-execution skip plan.
5//!
6//! When `nmbrs refine` re-attaches to an existing session, it
7//! reads the session's `phase_outcomes` table and builds a
8//! [`RefinePlan`] — the set of (phase_name, phase_labels)
9//! pairs that have already completed across any prior
10//! execution, plus the next execution id to record outcomes
11//! under.
12//!
13//! The plan rides on the executor context (`ExecCtx::refine_plan`)
14//! and the phase-walk gate checks each phase against
15//! [`is_completed`] before dispatching `run_phase`. Skipped
16//! phases still get their scene-tree node pushed (so the TUI /
17//! progress display shows them with a "skipped — prior outcome"
18//! status) but no cycles run and no new outcome row is written.
19//!
20//! Scope (MVP): `--scope=missing` only — skip phases whose
21//! exact identity already has a `completed` outcome row.
22//! `--scope=changed` (hash compare) and `--scope=all` are
23//! follow-up pushes. `--on-removed=` policies likewise deferred.
24
25use std::collections::HashSet;
26use std::path::Path;
27
28/// Pre-computed skip set + next-execution id for one refine
29/// invocation.
30#[derive(Debug, Clone)]
31pub struct RefinePlan {
32    /// `(phase_name, phase_labels)` pairs that have at least
33    /// one prior outcome with status `"completed"` across any
34    /// execution of the session. The phase-walk gate checks
35    /// each phase against this set before dispatching its
36    /// per-cycle work.
37    pub completed: HashSet<(String, String)>,
38    /// Every `(phase_name, phase_labels)` pair that has ANY
39    /// prior outcome (regardless of status). Used by the
40    /// `--on-removed=` policy to detect phases that exist in
41    /// the session's history but no longer appear in the
42    /// freshly pre-mapped workload — those are candidates
43    /// for the error / keep / drop decision.
44    pub seen_identities: HashSet<(String, String)>,
45    /// Prior `(name, labels) → provenance` for completed
46    /// phases: the BASE hash (`phase_outcomes.phase_hash`) plus
47    /// the SRD-107 consumed-params JSON. Used by the hash gates:
48    /// at phase activation the executor computes the current
49    /// base hash and param digests and compares via
50    /// [`Self::unchanged_verdict`]. Legacy rows (either field
51    /// NULL) always flag as changed, so the conservative
52    /// behavior runs the phase rather than wrongly skipping it.
53    pub completed_hashes: std::collections::HashMap<(String, String), PriorCompletion>,
54    /// The execution id this refine invocation will record
55    /// new outcomes under. One greater than the maximum
56    /// `exec_id` observed in the prior `phase_outcomes` rows;
57    /// at least `1` for sessions with no prior outcomes
58    /// (degenerate, but supported).
59    pub next_exec_id: u64,
60    /// Total prior outcome rows examined. Surfaced in the
61    /// startup log so the operator can sanity-check the
62    /// session their refine attached to.
63    pub prior_outcomes_seen: usize,
64    /// SRD-77 scope mode: `Missing` (skip prior-completed),
65    /// `Changed` (skip prior-completed AND prior_hash matches
66    /// current_hash), or `All` (no skip — empty `completed`).
67    /// Set by the runner from the `scope=` CLI param; the
68    /// executor's phase walk consults it to decide which gate
69    /// to apply.
70    pub scope: RefineScope,
71}
72
73/// The chronologically-latest completed outcome's provenance
74/// for one phase identity (SRD-77 base hash + SRD-107
75/// consumed-params JSON).
76#[derive(Debug, Clone, Default)]
77pub struct PriorCompletion {
78    pub phase_hash: Option<String>,
79    pub params_consumed: Option<String>,
80}
81
82/// Why a phase may NOT skip under the refine hash gate
83/// (SRD-107 Push 3) — surfaced in diagnostics so an operator
84/// sees "re-running load_train: param 'dataset' changed"
85/// instead of a bare hash mismatch.
86#[derive(Debug, Clone, PartialEq, Eq)]
87pub enum SkipBlocker {
88    /// No prior completed outcome carries comparable provenance
89    /// (new phase, legacy row, or unreadable stored map).
90    NoPrior,
91    /// The base hash differs: an enclosing scope program or the
92    /// phase's own declared config changed.
93    BaseChanged,
94    /// A consumed param's value changed (or the param is no
95    /// longer present). Carries the param name.
96    ParamChanged(String),
97}
98
99impl std::fmt::Display for SkipBlocker {
100    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
101        match self {
102            Self::NoPrior => write!(f, "no comparable prior outcome"),
103            Self::BaseChanged => write!(f, "scope or phase config changed"),
104            Self::ParamChanged(name) => write!(f, "param '{name}' changed"),
105        }
106    }
107}
108
109/// SRD-77 `--scope=` modes.
110#[derive(Debug, Clone, Copy, PartialEq, Eq)]
111pub enum RefineScope {
112    /// Skip every phase identity already completed in any
113    /// prior execution of this session. The default.
114    Missing,
115    /// Skip every phase identity whose `prior_hash ==
116    /// current_hash` (program shape unchanged). Phases with
117    /// no prior completion fall through; phases with prior
118    /// completion but a different hash re-run.
119    Changed,
120    /// Run every phase. Prior outcomes are preserved as
121    /// cardinal history under their original exec_id; the
122    /// new run writes under the bumped exec_id.
123    All,
124}
125
126/// SRD-77 — Every read-side path that touches session data is
127/// **execution-qualified**: it accepts an [`ExecutionQualifier`]
128/// at the call boundary, and the storage layer applies a
129/// matching `exec_id` filter to its queries. The "aggregate
130/// across every execution" intent is the explicit
131/// [`ExecutionQualifier::All`] variant, not an unqualified
132/// default — callers can never accidentally read across
133/// multiple executions when they meant the latest.
134///
135/// Construct via:
136/// - [`ExecutionQualifier::latest`] — resolves `max(exec_id)`
137///   from the session db at call time; the natural
138///   no-flag default for read commands.
139/// - [`ExecutionQualifier::specific(n)`] — single execution
140///   id (typically from a `--execution=<n>` CLI flag).
141/// - [`ExecutionQualifier::all`] — every execution; the
142///   `--all-executions` CLI flag.
143#[derive(Debug, Clone, Copy, PartialEq, Eq)]
144pub enum ExecutionQualifier {
145    /// One specific `exec_id`. The storage layer applies a
146    /// `WHERE exec_id = <n>` filter.
147    Specific(u64),
148    /// Every recorded execution. The storage layer applies
149    /// no `exec_id` filter — the "aggregate across all
150    /// executions" semantic, opted into explicitly.
151    All,
152}
153
154/// SRD-77 — the reserved CLI-side virtual qualifier that
155/// means "the most recent execution recorded in the session
156/// store". Resolvers translate this to a concrete exec_id at
157/// query-construction time. **Must never appear in stored
158/// data** — the metric_instance reserved-word guard refuses
159/// any write carrying `session="latest"` or
160/// `exec_id="latest"`.
161pub const LATEST_LITERAL: &str = "latest";
162
163/// SRD-77 — emit a per-command banner when the active session
164/// has more than one execution in its history. Tells the
165/// operator which execution the implicit `latest` default
166/// resolved to and lists the latest three in temporal order
167/// (newest first), with `<-- latest` marking the one their
168/// query is currently bound to.
169///
170/// Silent for sessions with 0 or 1 executions — the banner is
171/// only useful when an ambiguity actually exists.
172///
173/// Writes to stderr so it lands next to other operator-
174/// visible logging without colliding with stdout pipelines
175/// (a piped `nmbrs report ... | less` still gets clean stdout).
176pub fn warn_multi_execution_default(session_dir: &std::path::Path) {
177    let db_path = session_dir.join("metrics.db");
178    if !db_path.exists() {
179        return;
180    }
181    let conn = match rusqlite::Connection::open_with_flags(
182        &db_path,
183        rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
184    ) {
185        Ok(c) => c,
186        Err(_) => return,
187    };
188    let exists: bool = conn
189        .query_row(
190            "SELECT EXISTS(SELECT 1 FROM sqlite_master \
191         WHERE type='table' AND name='executions')",
192            [],
193            |r| r.get::<_, i64>(0),
194        )
195        .map(|n| n != 0)
196        .unwrap_or(false);
197    if !exists {
198        return;
199    }
200    let mut stmt = match conn.prepare(
201        "SELECT exec_id, verb, scope, disposition, started_at_nanos \
202         FROM executions ORDER BY exec_id DESC LIMIT 3",
203    ) {
204        Ok(s) => s,
205        Err(_) => return,
206    };
207    // One execution row: (exec_id, verb, scope, disposition, started_at_nanos).
208    type ExecRow = (i64, String, Option<String>, Option<String>, i64);
209    let rows: Vec<ExecRow> = match stmt.query_map([], |r| {
210        Ok((
211            r.get::<_, i64>(0)?,
212            r.get::<_, String>(1)?,
213            r.get::<_, Option<String>>(2)?,
214            r.get::<_, Option<String>>(3)?,
215            r.get::<_, i64>(4)?,
216        ))
217    }) {
218        Ok(it) => it.filter_map(Result::ok).collect(),
219        Err(_) => return,
220    };
221    let total: i64 = conn
222        .query_row("SELECT COUNT(*) FROM executions", [], |r| r.get(0))
223        .unwrap_or(0);
224    if total < 2 || rows.is_empty() {
225        return;
226    }
227    eprintln!(
228        "session has {total} execution(s); implicit qualifier `exec_id=latest` \
229         resolved to exec_id={latest}. recent (newest first):",
230        latest = rows[0].0,
231    );
232    for (i, (exec_id, verb, scope, disposition, _)) in rows.iter().enumerate() {
233        let marker = if i == 0 { " <-- latest" } else { "" };
234        let scope_part = scope
235            .as_deref()
236            .map(|s| format!(" scope={s}"))
237            .unwrap_or_default();
238        let disp_part = disposition
239            .as_deref()
240            .map(|d| format!(" {d}"))
241            .unwrap_or_else(|| " (in-flight)".to_string());
242        eprintln!("  exec_id={exec_id} verb={verb}{scope_part}{disp_part}{marker}");
243    }
244    eprintln!("  (pass `--execution=<n>` to target one, `--all-executions` to aggregate)");
245}
246
247impl ExecutionQualifier {
248    /// Single execution id.
249    pub fn specific(n: u64) -> Self {
250        Self::Specific(n)
251    }
252
253    /// Aggregate across every execution.
254    pub fn all() -> Self {
255        Self::All
256    }
257
258    /// Resolve "the most recent execution" against the
259    /// session db at `session_dir`. Returns
260    /// [`Self::Specific(max_exec_id)`] when at least one
261    /// execution is recorded, falling back to
262    /// [`Self::Specific(1)`] when the db is empty / absent
263    /// (so the qualifier still narrows to a specific id —
264    /// the caller still gets the explicit qualification
265    /// promise, just against an empty target).
266    pub fn latest(session_dir: &std::path::Path) -> Self {
267        match latest_exec_id_for_session(session_dir) {
268            Some(n) => Self::Specific(n),
269            None => Self::Specific(1),
270        }
271    }
272
273    /// True iff this qualifier matches every recorded
274    /// execution. Storage-layer query builders use this to
275    /// decide whether to attach the `WHERE exec_id = …`
276    /// clause.
277    pub fn matches_all(&self) -> bool {
278        matches!(self, Self::All)
279    }
280
281    /// The specific `exec_id` when narrowed, otherwise
282    /// `None`. Storage-layer builders use this to bind the
283    /// `WHERE exec_id = ?` parameter.
284    pub fn specific_id(&self) -> Option<u64> {
285        match self {
286            Self::Specific(n) => Some(*n),
287            Self::All => None,
288        }
289    }
290}
291
292/// Read the maximum `exec_id` recorded in the session's
293/// `phase_outcomes` table — i.e. "which execution_id is the
294/// most recent one in this session's history". Used by
295/// [`ExecutionQualifier::latest`] to resolve the latest-
296/// execution intent into a concrete id.
297///
298/// Returns `None` when:
299/// - The session dir doesn't exist
300/// - The sqlite file doesn't exist (no run captured yet)
301/// - The `phase_outcomes` table is empty or absent
302///
303/// `O(1)` against the PK index — no full scan.
304pub fn latest_exec_id_for_session(session_dir: &std::path::Path) -> Option<u64> {
305    let db_path = session_dir.join("metrics.db");
306    if !db_path.exists() {
307        return None;
308    }
309    let conn = rusqlite::Connection::open_with_flags(
310        &db_path,
311        rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
312    )
313    .ok()?;
314    let table_exists: bool = conn
315        .query_row(
316            "SELECT EXISTS(SELECT 1 FROM sqlite_master \
317         WHERE type='table' AND name='phase_outcomes')",
318            [],
319            |r| r.get::<_, i64>(0),
320        )
321        .map(|n| n != 0)
322        .unwrap_or(false);
323    if !table_exists {
324        return None;
325    }
326    let max: Option<i64> = conn
327        .query_row("SELECT MAX(exec_id) FROM phase_outcomes", [], |r| {
328            r.get::<_, Option<i64>>(0)
329        })
330        .ok()
331        .flatten();
332    max.map(|v| v.max(0) as u64)
333}
334
335impl RefinePlan {
336    /// `O(1)` check used by the executor's phase-walk gate.
337    pub fn is_completed(&self, phase_name: &str, phase_labels: &str) -> bool {
338        self.completed
339            .contains(&(phase_name.to_string(), phase_labels.to_string()))
340    }
341
342    /// SRD-77 `--scope=changed` — true iff this phase has a
343    /// prior completed outcome AND the prior outcome's
344    /// `phase_hash` matches `current_hex`. False when:
345    /// - No prior completion: phase is new → run.
346    /// - Prior completion but hash differs: program shape
347    ///   changed → re-run.
348    /// - Prior completion with `phase_hash = NULL`: legacy
349    ///   row from before the column was added → conservatively
350    ///   re-run (we can't prove unchanged, so default to "do
351    ///   the work").
352    pub fn is_unchanged(
353        &self,
354        phase_name: &str,
355        phase_labels: &str,
356        current_hex: &str,
357        current_params: &std::collections::HashMap<String, String>,
358    ) -> bool {
359        self.unchanged_verdict(phase_name, phase_labels, current_hex, current_params)
360            .is_ok()
361    }
362
363    /// SRD-107 Push 3 — the three-way skip-validity check with a
364    /// NAMED blocker on failure: base hash equal AND every stored
365    /// consumed param's current value digests to its stored
366    /// digest. Conservative on any gap (legacy rows, unreadable
367    /// stored map): re-run rather than wrongly skip.
368    pub fn unchanged_verdict(
369        &self,
370        phase_name: &str,
371        phase_labels: &str,
372        current_hex: &str,
373        current_params: &std::collections::HashMap<String, String>,
374    ) -> Result<(), SkipBlocker> {
375        let key = (phase_name.to_string(), phase_labels.to_string());
376        let Some(prior) = self.completed_hashes.get(&key) else {
377            return Err(SkipBlocker::NoPrior);
378        };
379        match prior.phase_hash.as_deref() {
380            None => return Err(SkipBlocker::NoPrior),
381            Some(prior_hex) if prior_hex != current_hex => return Err(SkipBlocker::BaseChanged),
382            Some(_) => {}
383        }
384        let Some(json) = prior.params_consumed.as_deref() else {
385            // A base-matching row without the SRD-107 map should
386            // not exist post-upgrade; treat as incomparable.
387            return Err(SkipBlocker::NoPrior);
388        };
389        let Ok(stored) = serde_json::from_str::<std::collections::BTreeMap<String, String>>(json)
390        else {
391            return Err(SkipBlocker::NoPrior);
392        };
393        for (name, stored_digest) in stored {
394            let current = current_params
395                .get(&name)
396                .map(|v| crate::checkpoint::params_scope::value_digest(v));
397            if current.as_deref() != Some(stored_digest.as_str()) {
398                return Err(SkipBlocker::ParamChanged(name));
399            }
400        }
401        Ok(())
402    }
403
404    /// Should this phase be skipped per the plan's scope?
405    /// Centralises the scope→gate dispatch so the executor
406    /// walker stays a single conditional. The hash arg is
407    /// only consulted under `Changed`; passing `""` is fine
408    /// for `Missing` / `All` callers that don't know it yet.
409    pub fn should_skip(
410        &self,
411        phase_name: &str,
412        phase_labels: &str,
413        current_hex: &str,
414        current_params: &std::collections::HashMap<String, String>,
415    ) -> bool {
416        match self.scope {
417            RefineScope::All => false,
418            RefineScope::Missing => self.is_completed(phase_name, phase_labels),
419            RefineScope::Changed => {
420                self.is_unchanged(phase_name, phase_labels, current_hex, current_params)
421            }
422        }
423    }
424
425    /// Open the session directory's `metrics.db`, read every
426    /// `phase_outcomes` row, and compute the skip plan.
427    ///
428    /// Returns `None` when:
429    /// - The session directory doesn't exist
430    /// - The sqlite file doesn't exist (no prior run captured outcomes)
431    /// - The sqlite file exists but the `phase_outcomes` table is missing
432    ///   (legacy session predating SRD-76)
433    ///
434    /// Sqlite errors mid-query log at WARN and produce an
435    /// empty plan — refine falls back to "run everything" rather
436    /// than failing the invocation, on the principle that a
437    /// half-readable database shouldn't block the operator
438    /// from making progress.
439    pub fn load_from_session_dir(session_dir: &Path) -> Option<Self> {
440        let db_path = session_dir.join("metrics.db");
441        if !db_path.exists() {
442            return None;
443        }
444        let conn = match rusqlite::Connection::open_with_flags(
445            &db_path,
446            rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
447        ) {
448            Ok(c) => c,
449            Err(e) => {
450                crate::diag!(
451                    crate::observer::LogLevel::Warn,
452                    "refine: failed to open {}: {e}",
453                    db_path.display()
454                );
455                return None;
456            }
457        };
458        // Verify the table exists before querying — a session
459        // dir from before SRD-76 lands won't have it, and
460        // `prepare` against a missing table errors with a
461        // message that's noisier than "no plan available".
462        let exists: bool = conn
463            .query_row(
464                "SELECT EXISTS(SELECT 1 FROM sqlite_master \
465             WHERE type='table' AND name='phase_outcomes')",
466                [],
467                |r| r.get::<_, i64>(0),
468            )
469            .map(|n| n != 0)
470            .unwrap_or(false);
471        if !exists {
472            return None;
473        }
474        // SRD-107 legacy-read guard: the params_consumed column
475        // may be absent on dbs never re-opened by a current
476        // writer (this connection is read-only, so no migration
477        // here); an absent column reads as NULL.
478        let has_params_col: bool = conn
479            .prepare("PRAGMA table_info(phase_outcomes)")
480            .ok()
481            .and_then(|mut s| {
482                let mut found = false;
483                let mut rows = s.query([]).ok()?;
484                while let Ok(Some(r)) = rows.next() {
485                    if r.get::<_, String>(1)
486                        .map(|n| n == "params_consumed")
487                        .unwrap_or(false)
488                    {
489                        found = true;
490                    }
491                }
492                Some(found)
493            })
494            .unwrap_or(false);
495        let pc_col = if has_params_col {
496            "params_consumed"
497        } else {
498            "NULL"
499        };
500        let mut stmt = match conn.prepare(&format!(
501            "SELECT exec_id, phase_name, phase_labels, status, phase_hash, \
502                    {pc_col}, ended_at_nanos \
503             FROM phase_outcomes \
504             ORDER BY ended_at_nanos"
505        )) {
506            Ok(s) => s,
507            Err(e) => {
508                crate::diag!(
509                    crate::observer::LogLevel::Warn,
510                    "refine: failed to prepare query: {e}"
511                );
512                return None;
513            }
514        };
515        let rows = match stmt.query_map([], |row| {
516            Ok((
517                row.get::<_, i64>(0)?,
518                row.get::<_, String>(1)?,
519                row.get::<_, String>(2)?,
520                row.get::<_, String>(3)?,
521                row.get::<_, Option<String>>(4)?,
522                row.get::<_, Option<String>>(5)?,
523            ))
524        }) {
525            Ok(r) => r,
526            Err(e) => {
527                crate::diag!(
528                    crate::observer::LogLevel::Warn,
529                    "refine: failed to query phase_outcomes: {e}"
530                );
531                return None;
532            }
533        };
534        let mut completed = HashSet::new();
535        let mut seen_identities: HashSet<(String, String)> = HashSet::new();
536        let mut completed_hashes: std::collections::HashMap<(String, String), PriorCompletion> =
537            std::collections::HashMap::new();
538        let mut max_exec_id: u64 = 0;
539        let mut count: usize = 0;
540        // Rows arrive ordered by `ended_at_nanos`, so the
541        // chronologically latest completed outcome's hash wins
542        // for a given (name, labels). This is what we want for
543        // `scope=changed`: "did the LAST completed run match
544        // what we'd compute now?"
545        for row in rows.flatten() {
546            let (exec_id, name, labels, status, phase_hash, params_consumed) = row;
547            let exec_id = exec_id.max(0) as u64;
548            if exec_id > max_exec_id {
549                max_exec_id = exec_id;
550            }
551            count += 1;
552            seen_identities.insert((name.clone(), labels.clone()));
553            if status == "completed" {
554                completed.insert((name.clone(), labels.clone()));
555                completed_hashes.insert(
556                    (name, labels),
557                    PriorCompletion {
558                        phase_hash,
559                        params_consumed,
560                    },
561                );
562            }
563        }
564        // SRD-77 — `next_exec_id` must consult the `executions`
565        // table too. A prior refine that skipped every phase
566        // writes ZERO phase_outcomes rows but DOES insert an
567        // executions row, so phase_outcomes alone would miss
568        // the bump and the next invocation would collide on the
569        // executions PK. Take MAX across both sources.
570        let executions_max: u64 = conn
571            .query_row(
572                "SELECT EXISTS(SELECT 1 FROM sqlite_master \
573             WHERE type='table' AND name='executions')",
574                [],
575                |r| r.get::<_, i64>(0),
576            )
577            .map(|n| n != 0)
578            .ok()
579            .filter(|exists| *exists)
580            .and_then(|_| {
581                conn.query_row("SELECT MAX(exec_id) FROM executions", [], |r| {
582                    r.get::<_, Option<i64>>(0)
583                })
584                .ok()
585                .flatten()
586            })
587            .map(|v| v.max(0) as u64)
588            .unwrap_or(0);
589        let max_exec_id = max_exec_id.max(executions_max);
590        Some(Self {
591            completed,
592            seen_identities,
593            completed_hashes,
594            next_exec_id: max_exec_id + 1,
595            prior_outcomes_seen: count,
596            scope: RefineScope::Missing,
597        })
598    }
599}
600
601#[cfg(test)]
602mod tests {
603    use super::*;
604
605    #[test]
606    fn missing_db_returns_none() {
607        let tmp = tempfile::tempdir().unwrap();
608        assert!(RefinePlan::load_from_session_dir(tmp.path()).is_none());
609    }
610
611    #[test]
612    fn db_without_phase_outcomes_returns_none() {
613        let tmp = tempfile::tempdir().unwrap();
614        let db_path = tmp.path().join("metrics.db");
615        let conn = rusqlite::Connection::open(&db_path).unwrap();
616        conn.execute("CREATE TABLE other (x INTEGER)", []).unwrap();
617        drop(conn);
618        assert!(RefinePlan::load_from_session_dir(tmp.path()).is_none());
619    }
620
621    fn make_db_with_outcomes(dir: &Path, rows: &[(u64, &str, &str, &str)]) {
622        let db_path = dir.join("metrics.db");
623        let conn = rusqlite::Connection::open(&db_path).unwrap();
624        conn.execute(
625            "CREATE TABLE phase_outcomes (
626                session       TEXT    NOT NULL,
627                exec_id       INTEGER NOT NULL,
628                phase_name    TEXT    NOT NULL,
629                phase_labels  TEXT    NOT NULL,
630                status        TEXT    NOT NULL,
631                duration_secs REAL    NOT NULL DEFAULT 0,
632                started_at_nanos INTEGER NOT NULL DEFAULT 0,
633                ended_at_nanos   INTEGER NOT NULL DEFAULT 0,
634                phase_hash    TEXT,
635                PRIMARY KEY (session, exec_id, phase_name, phase_labels)
636            )",
637            [],
638        )
639        .unwrap();
640        for (exec, name, labels, status) in rows {
641            conn.execute(
642                "INSERT INTO phase_outcomes (session, exec_id, phase_name, phase_labels, status) \
643                 VALUES ('s', ?1, ?2, ?3, ?4)",
644                rusqlite::params![*exec as i64, name, labels, status],
645            )
646            .unwrap();
647        }
648    }
649
650    #[test]
651    fn completed_phases_populate_skip_set() {
652        let tmp = tempfile::tempdir().unwrap();
653        make_db_with_outcomes(
654            tmp.path(),
655            &[
656                (1, "schema", "", "completed"),
657                (1, "load_data", "k=10", "completed"),
658                (1, "query", "k=10,limit=20", "failed"),
659            ],
660        );
661        let plan = RefinePlan::load_from_session_dir(tmp.path()).unwrap();
662        assert_eq!(plan.next_exec_id, 2);
663        assert_eq!(plan.prior_outcomes_seen, 3);
664        assert!(plan.is_completed("schema", ""));
665        assert!(plan.is_completed("load_data", "k=10"));
666        // failed phases must NOT be in the skip set — refine
667        // should re-run them.
668        assert!(!plan.is_completed("query", "k=10,limit=20"));
669    }
670
671    #[test]
672    fn next_exec_id_bumps_past_max_prior() {
673        let tmp = tempfile::tempdir().unwrap();
674        make_db_with_outcomes(
675            tmp.path(),
676            &[
677                (1, "schema", "", "completed"),
678                (3, "query", "", "completed"),
679                (2, "load", "", "completed"),
680            ],
681        );
682        let plan = RefinePlan::load_from_session_dir(tmp.path()).unwrap();
683        assert_eq!(plan.next_exec_id, 4);
684    }
685
686    #[test]
687    fn empty_table_yields_exec_id_1() {
688        let tmp = tempfile::tempdir().unwrap();
689        make_db_with_outcomes(tmp.path(), &[]);
690        let plan = RefinePlan::load_from_session_dir(tmp.path()).unwrap();
691        assert_eq!(plan.next_exec_id, 1);
692        assert_eq!(plan.prior_outcomes_seen, 0);
693        assert!(plan.completed.is_empty());
694    }
695
696    // ── ExecutionQualifier ────────────────────────────────
697
698    #[test]
699    fn execution_qualifier_specific_carries_id() {
700        let q = ExecutionQualifier::specific(7);
701        assert_eq!(q.specific_id(), Some(7));
702        assert!(!q.matches_all());
703    }
704
705    #[test]
706    fn execution_qualifier_all_carries_no_id() {
707        let q = ExecutionQualifier::all();
708        assert_eq!(q.specific_id(), None);
709        assert!(q.matches_all());
710    }
711
712    /// `latest()` resolves to `Specific(max_exec_id)` when the
713    /// session db carries phase_outcomes — pins the "no
714    /// implicit aggregation default" invariant: even the
715    /// `latest` constructor narrows to one execution.
716    #[test]
717    fn execution_qualifier_latest_resolves_to_max_exec_id() {
718        let tmp = tempfile::tempdir().unwrap();
719        make_db_with_outcomes(
720            tmp.path(),
721            &[
722                (1, "p1", "", "completed"),
723                (3, "p2", "", "completed"),
724                (2, "p3", "", "completed"),
725            ],
726        );
727        let q = ExecutionQualifier::latest(tmp.path());
728        assert_eq!(
729            q.specific_id(),
730            Some(3),
731            "latest MUST resolve to max(exec_id)=3"
732        );
733        assert!(
734            !q.matches_all(),
735            "latest MUST narrow to a specific id, not aggregate"
736        );
737    }
738
739    /// Empty db falls back to `Specific(1)` rather than
740    /// silently degrading to aggregate. This is the
741    /// "qualifier always narrows" promise: read paths can rely
742    /// on a concrete exec_id even on a session with no prior
743    /// runs.
744    #[test]
745    fn execution_qualifier_latest_on_empty_db_yields_specific_1() {
746        let tmp = tempfile::tempdir().unwrap();
747        let q = ExecutionQualifier::latest(tmp.path());
748        assert_eq!(
749            q.specific_id(),
750            Some(1),
751            "latest on empty db MUST yield Specific(1), not All"
752        );
753    }
754
755    /// Db with no `phase_outcomes` table at all (legacy /
756    /// pre-SRD-77 session) must STILL produce Specific(1) —
757    /// callers can't drift into aggregate just because the
758    /// table is missing.
759    #[test]
760    fn execution_qualifier_latest_on_missing_table_yields_specific_1() {
761        let tmp = tempfile::tempdir().unwrap();
762        let db_path = tmp.path().join("metrics.db");
763        let conn = rusqlite::Connection::open(&db_path).unwrap();
764        conn.execute("CREATE TABLE other (x INTEGER)", []).unwrap();
765        drop(conn);
766        let q = ExecutionQualifier::latest(tmp.path());
767        assert_eq!(q.specific_id(), Some(1));
768    }
769}