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}