iris: the lint over a portal site, as its own repo
Test / test (push) Failing after 12s
Publish release / publish (push) Failing after 4s

question_lint becomes `iris check`, needs_replay becomes `iris
replay`, and the needs module and local-checkout loaders come with
them. Portal is a library dependency pinned to the release iris
matches (v0.3.36), so every type iris reads is the site's own and the
two never disagree about what a page is. `iris --path questions`
still works, so a content repo's one-line CI needs only a new path.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01L4jrCgLiKKHAEFuZUJjckH
This commit is contained in:
Bendik Aagaard Lynghaug
2026-09-22 17:25:52 +02:00
co-authored by Claude Fable 5.1
commit 73906cdf35
12 changed files with 6927 additions and 0 deletions
+225
View File
@@ -0,0 +1,225 @@
//! Headless schema-check binary - loads content and `aggregates.yaml`
//! (from a Gitea repo URL or a local directory) and validates them
//! exactly the way `content::watch_for_reload`/`main.rs`'s boot path
//! do, with no NATS, OIDC, web server, or JetStream connection
//! involved. Built once by portal's own deploy workflow and run
//! directly by `questions`' own CI (same bare-metal runner/host,
//! published to a stable path - no artifact download needed), rather
//! than compiled there - keeps that repo's CI coupling to "run a
//! static binary," not "build a Rust workspace."
use crate::local::{load_aggregates_from_dir, load_from_dir, load_needs_from_dir, load_site_from_dir};
use crate::needs;
use portal::{aggregates, content};
pub async fn run(argv: Vec<String>) -> anyhow::Result<()> {
let mut args = argv.into_iter();
let mut repo: Option<String> = None;
let mut branch = "main".to_string();
let mut subdir = "questions".to_string();
let mut path: Option<String> = None;
let mut tasks_out: Option<String> = None;
let mut answers_in: Option<String> = None;
let mut min_accuracy: f64 = 0.0;
let mut skim = false;
let mut sim_out: Option<String> = None;
while let Some(arg) = args.next() {
match arg.as_str() {
"--repo" => repo = args.next(),
"--branch" => branch = args.next().unwrap_or(branch),
"--subdir" => subdir = args.next().unwrap_or(subdir),
"--path" => path = args.next(),
"--needs-tasks" => tasks_out = args.next(),
"--skim" => skim = true,
"--needs-sim" => sim_out = args.next(),
"--needs-score" => answers_in = args.next(),
"--min-accuracy" => {
min_accuracy = args.next().and_then(|v| v.parse().ok()).unwrap_or(min_accuracy)
}
other => {
eprintln!("unknown argument: {other}");
std::process::exit(2);
}
}
}
let (questions, aggregates_map, site, needs_raw) = if let Some(dir) = path {
(
load_from_dir(&dir)?,
load_aggregates_from_dir(&dir)?,
load_site_from_dir(&dir)?,
load_needs_from_dir(&dir)?,
)
} else if let Some(repo_url) = repo {
let questions = content::load_questions_from_gitea(&repo_url, &branch, &subdir).await?;
let aggregates_map = content::load_aggregates_from_gitea(&repo_url, &branch).await?;
// load_site_from_gitea already validates; a missing file is
// the default config, same as at portal boot.
let site = content::load_site_from_gitea(&repo_url, &branch).await?;
let needs_raw = content::load_needs_from_gitea(&repo_url, &branch).await?;
(questions, aggregates_map, site, needs_raw)
} else {
eprintln!(
"usage: iris check (--repo <gitea-url> [--branch main] [--subdir questions] | --path <local-dir>) [--needs-tasks <out.jsonl> [--skim]] [--needs-sim <model.json>] [--needs-score <answers.jsonl> [--min-accuracy 0.8]]"
);
std::process::exit(2);
};
if let Err(e) = site.validate() {
eprintln!("FAIL: {e}");
std::process::exit(1);
}
let aggregates_map = content::with_builtin_aggregates(aggregates_map);
match content::validate_questions(&questions, &aggregates_map) {
Ok(()) => {
// The runtime's own event desk: a repo announcing pages
// (Question.event) should read portal_events somewhere, or
// "post what happened" tasks pile up unseen. Warn, don't
// fail - the runtime creates the records either way.
if questions.values().any(|q| q.event.is_some())
&& !content::bucket_is_read(&questions, content::EVENTS_BUCKET)
{
eprintln!(
"WARN: pages carry `event:` but no kv resource reads bucket {:?} - add a desk so summaries get posted",
content::EVENTS_BUCKET
);
}
println!(
"OK: {} question(s), {} aggregate(s) valid",
questions.len(),
aggregates_map.len()
);
check_needs(
needs_raw.as_deref(),
&questions,
&aggregates_map,
tasks_out.as_deref(),
skim,
sim_out.as_deref(),
answers_in.as_deref(),
min_accuracy,
)
}
Err(e) => {
eprintln!("FAIL: {e}");
std::process::exit(1);
}
}
}
/// The business-needs pass (see `portal::needs`): structure always,
/// then the wording tasks written out and/or an engine's answers
/// scored, when asked. A repo without `needs.yaml` skips all of it.
fn check_needs(
raw: Option<&str>,
questions: &std::collections::HashMap<String, content::Question>,
aggregates_map: &std::collections::HashMap<String, aggregates::AggregateSchema>,
tasks_out: Option<&str>,
skim: bool,
sim_out: Option<&str>,
answers_in: Option<&str>,
min_accuracy: f64,
) -> anyhow::Result<()> {
let Some(raw) = raw else {
if tasks_out.is_some() || answers_in.is_some() {
eprintln!("FAIL: --needs-tasks/--needs-score given, but the repo has no needs.yaml");
std::process::exit(1);
}
return Ok(());
};
let file = match needs::parse(raw) {
Ok(file) => file,
Err(e) => {
eprintln!("FAIL: {e}");
std::process::exit(1);
}
};
let findings = needs::check(&file, questions, aggregates_map);
let mut failed = 0;
for finding in &findings {
let level = match finding.level {
needs::Level::Fail => {
failed += 1;
"FAIL"
}
needs::Level::Warn => "WARN",
};
eprintln!("{level}: need {:?}: {}", finding.need, finding.message);
}
for need in &file.needs {
if let Some(path) = needs::path_for(questions, need) {
let steps: Vec<String> = path
.iter()
.map(|s| format!("{} [{}]", s.page, s.alternative))
.collect();
println!("need {:?}: {} -> {:?}", need.id, steps.join(" -> "), need.lands_in);
}
}
if failed == 0 {
let planned = file.needs.iter().filter(|n| n.planned).count();
println!(
"OK: {} need(s) structurally met, {} planned, {} persona(s)",
file.needs.len() - planned,
planned,
file.personas.len()
);
}
// The wording pass still runs for the needs that do have a path -
// a closed flow elsewhere is no reason to stop grading the rest.
if let Some(out) = sim_out {
let model = needs::sim_model(&file, questions, aggregates_map);
std::fs::write(out, serde_json::to_string_pretty(&model)?)
.map_err(|e| anyhow::anyhow!("writing {out}: {e}"))?;
println!("wrote the simulation model to {out}");
}
let tasks = needs::tasks(&file, questions, skim);
if let Some(out) = tasks_out {
let mut lines = String::new();
for task in &tasks {
lines.push_str(&serde_json::to_string(task)?);
lines.push('\n');
}
std::fs::write(out, lines).map_err(|e| anyhow::anyhow!("writing {out}: {e}"))?;
println!("wrote {} wording task(s) to {out}", tasks.len());
}
if let Some(path) = answers_in {
let raw = std::fs::read_to_string(path).map_err(|e| anyhow::anyhow!("reading {path}: {e}"))?;
let mut answers = Vec::new();
for (n, line) in raw.lines().enumerate().filter(|(_, l)| !l.trim().is_empty()) {
answers.push(
serde_json::from_str::<needs::Answer>(line)
.map_err(|e| anyhow::anyhow!("{path}:{}: {e}", n + 1))?,
);
}
let score = needs::score(&tasks, &answers, 0.6);
for (id, picked, expected, confidence) in &score.misroutes {
eprintln!(
"MISROUTE: {id}: picked {picked:?}{}, the need is met by {expected:?}",
confidence.map(|c| format!(" ({c:.2})")).unwrap_or_default()
);
}
for (id, confidence) in &score.hesitant {
eprintln!("HESITANT: {id}: right, at {confidence:.2}");
}
for id in &score.unanswered {
eprintln!("UNANSWERED: {id}");
}
println!(
"wording: {}/{} routed right ({:.0}%)",
score.correct,
score.total,
score.accuracy() * 100.0
);
if score.accuracy() < min_accuracy {
eprintln!("FAIL: below --min-accuracy {min_accuracy}");
std::process::exit(1);
}
}
if failed > 0 {
eprintln!("FAIL: {failed} unmet business need finding(s)");
std::process::exit(1);
}
Ok(())
}
+99
View File
@@ -0,0 +1,99 @@
//! Loading a content repo from a local checkout - the offline
//! counterparts to the Gitea loaders in `content`, for tools that lint
//! or replay a branch that has not been pushed (question_lint,
//! needs_replay). `dir` is always the pages directory (`questions/`);
//! `site.yaml`, `aggregates.yaml` and `needs.yaml` sit one level up.
use portal::{aggregates, content};
/// The offline counterpart to `content::load_site_from_gitea` -
/// `site.yaml` lives at the repo root like `aggregates.yaml`, and is
/// just as optional locally: missing means default branding.
pub fn load_site_from_dir(dir: &str) -> anyhow::Result<content::SiteConfig> {
let site_path = std::path::Path::new(dir)
.parent()
.unwrap_or_else(|| std::path::Path::new("."))
.join("site.yaml");
if !site_path.exists() {
return Ok(content::SiteConfig::default());
}
let raw = std::fs::read_to_string(&site_path)
.map_err(|e| anyhow::anyhow!("reading {}: {e}", site_path.display()))?;
serde_yaml::from_str(&raw).map_err(|e| anyhow::anyhow!("parsing {}: {e}", site_path.display()))
}
/// The offline counterpart to `content::load_aggregates_from_gitea` -
/// `aggregates.yaml` lives at the repo root, one level up from the
/// pages directory `--path` names, so `dir`'s parent is where it's
/// looked for. Missing entirely is not an error here (unlike
/// `--repo` mode, where `aggregates.yaml` is required) - a local
/// checkout being linted may not have one, and every declared
/// transition still gets checked against whatever *is* found; an
/// empty map just means nothing is checked.
pub fn load_aggregates_from_dir(
dir: &str,
) -> anyhow::Result<std::collections::HashMap<String, aggregates::AggregateSchema>> {
let aggregates_path = std::path::Path::new(dir)
.parent()
.unwrap_or_else(|| std::path::Path::new("."))
.join("aggregates.yaml");
if !aggregates_path.exists() {
return Ok(std::collections::HashMap::new());
}
let raw = std::fs::read_to_string(&aggregates_path)
.map_err(|e| anyhow::anyhow!("reading {}: {e}", aggregates_path.display()))?;
aggregates::parse_aggregates_yaml(&raw)
}
/// The offline counterpart to `content::load_questions_from_gitea` -
/// the same recursive tree walk and `content::build_questions`
/// pipeline (derived ids, relative refs, sections, followup
/// inference), just reading a local checkout instead of Gitea's API,
/// for linting a branch that hasn't been pushed yet.
pub fn load_from_dir(
dir: &str,
) -> anyhow::Result<std::collections::HashMap<String, content::Question>> {
let base = std::path::Path::new(dir);
let mut files = Vec::new();
collect_yaml(base, base, &mut files)?;
content::build_questions(&files)
}
fn collect_yaml(
base: &std::path::Path,
dir: &std::path::Path,
out: &mut Vec<(String, String)>,
) -> anyhow::Result<()> {
for entry in std::fs::read_dir(dir)? {
let path = entry?.path();
if path.is_dir() {
collect_yaml(base, &path, out)?;
} else if path.extension().and_then(|e| e.to_str()) == Some("yaml") {
let rel = path
.strip_prefix(base)
.expect("walked paths sit under their base")
.to_string_lossy()
.replace('\\', "/");
let raw = std::fs::read_to_string(&path)
.map_err(|e| anyhow::anyhow!("reading {}: {e}", path.display()))?;
out.push((rel, raw));
}
}
Ok(())
}
/// The offline counterpart to `content::load_needs_from_gitea` -
/// `needs.yaml` sits at the repo root beside `aggregates.yaml`.
pub fn load_needs_from_dir(dir: &str) -> anyhow::Result<Option<String>> {
let needs_path = std::path::Path::new(dir)
.parent()
.unwrap_or_else(|| std::path::Path::new("."))
.join("needs.yaml");
if !needs_path.exists() {
return Ok(None);
}
std::fs::read_to_string(&needs_path)
.map(Some)
.map_err(|e| anyhow::anyhow!("reading {}: {e}", needs_path.display()))
}
+56
View File
@@ -0,0 +1,56 @@
//! iris - the shield over a portal site. A content push is an incoming
//! wormhole; nothing comes through until it checks out.
//!
//! `iris check` loads a content repo (a Gitea URL or a local checkout)
//! exactly the way the running site does, validates it, and holds it to
//! the business needs it declares in `needs.yaml` - a path within
//! budget for every kind of visitor, a desk for every handling group,
//! every finished state reachable, no record stranding. It also exports
//! what a trainer needs to grade the wording: personas as typed choice
//! tasks, and a resolved simulation model.
//!
//! `iris replay` runs simulated cases through the site's real aggregate
//! engine on a JetStream and reports every bucket by state; with
//! `--report-only` it takes the same scorecard from a live site's
//! buckets.
//!
//! Every type iris reads is portal's own, pinned to the portal release
//! it matches, so the lint and the site never disagree about what a
//! page is.
mod check;
mod local;
mod needs;
mod replay;
#[tokio::main]
async fn main() -> anyhow::Result<()> {
let mut argv: Vec<String> = std::env::args().skip(1).collect();
// `iris --path questions` reads as `iris check --path questions`,
// so the content repos' one-line CI keeps working as it was.
match argv.first().map(String::as_str) {
Some("check") => {
argv.remove(0);
check::run(argv).await
}
Some("replay") => {
argv.remove(0);
replay::run(argv).await
}
Some("--version") | Some("-V") => {
println!("iris {} (portal {})", env!("CARGO_PKG_VERSION"), PORTAL_TAG);
Ok(())
}
Some(flag) if flag.starts_with("--") => check::run(argv).await,
_ => {
eprintln!(
"iris {} - matches portal {PORTAL_TAG}\n\nusage:\n iris check (--repo <gitea-url> [--branch main] [--subdir questions] | --path <dir>) [--needs-tasks <out.jsonl> [--skim]] [--needs-sim <model.json>] [--needs-score <answers.jsonl> [--min-accuracy 0.8]]\n iris replay --path <dir> --nats <url> (--cases <cases.jsonl> | --report-only)",
env!("CARGO_PKG_VERSION")
);
std::process::exit(2)
}
}
}
/// The portal release this iris was built against (see Cargo.toml).
const PORTAL_TAG: &str = "v0.3.36";
+968
View File
@@ -0,0 +1,968 @@
//! A business's needs, stated apart from how its pages word them.
//!
//! A content repo may ship a `needs.yaml` at its root: each *need* says
//! who wants what, where they enter, which bucket the answer must land
//! in, and which group has to be able to carry it to a finished state.
//! Two things are checked against it:
//!
//! - **Structure** (`check`): deterministic. A path of submissions
//! leads from the entry page to the bucket, within an effort budget;
//! the handling group has a desk over that bucket; every finished
//! state is reachable through the moves that desk offers, and no
//! record strands on the way.
//! - **Wording** (`tasks` / `score`): each *persona* is a visitor in
//! their own words. Every page on their path becomes one typed
//! `choice` question for an external answering engine - the page's
//! question, its alternatives as the options - and the engine's picks
//! are scored against the alternatives that actually meet the need.
//! A misroute is a wording defect: the page said one thing and the
//! visitor heard another.
//!
//! The engine is deliberately outside this crate: tasks go out as
//! JSONL, answers come back as JSONL, so any classifier or LLM can sit
//! in between and be compared on the same tasks.
use std::collections::{BTreeMap, HashMap, HashSet, VecDeque};
use serde::{Deserialize, Serialize};
use portal::aggregates::AggregateSchema;
use portal::content::{action_target_id, Alternative, Question};
#[derive(Debug, Default, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct NeedsFile {
#[serde(default)]
pub needs: Vec<Need>,
#[serde(default)]
pub personas: Vec<Persona>,
}
#[derive(Debug, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct Need {
pub id: String,
/// The kind of visitor - prose, for the reader of the report.
pub who: String,
pub wants: String,
/// The page they arrive on.
#[serde(default = "root")]
pub enters: String,
/// The group they hold when they arrive; unset means anonymous.
#[serde(default)]
pub enters_as: Option<String>,
/// The alternative on the entry page that serves this need, when
/// several record into the same bucket.
#[serde(default)]
pub via: Option<String>,
/// The bucket their answer must be recorded into.
pub lands_in: String,
/// The group that must be able to see the record and move it.
pub handled_by: String,
/// The states that count as handled.
pub done_when: Vec<String>,
/// Submissions the visitor may be asked for before the record
/// exists.
#[serde(default = "default_max_steps")]
pub max_steps: usize,
/// Required fields across the whole path - the onboarding effort
/// budget.
#[serde(default)]
pub max_fields: Option<usize>,
/// Declared, and knowingly not open yet: what is missing reports
/// as a warning, so the gap stays visible without failing CI.
#[serde(default)]
pub planned: bool,
/// The finished states that are a good outcome for the business -
/// a subset of `done_when`. "lost" is finished; "won" is success.
/// What a simulation or a live bucket report counts toward.
#[serde(default)]
pub success: Vec<String>,
/// Every stage the business said a case goes through, in its own
/// order, first to last. Written from how the business describes
/// its work, before any state machine exists, so that a machine
/// which quietly skips a stage ("we go and look, then we quote"
/// becoming open -> quoted) fails instead of passing.
#[serde(default)]
pub stages: Vec<String>,
/// Other groups whose desks carry some of this need's moves: a
/// site manager closing what a project manager opened. Their
/// buttons count toward reaching `done_when` exactly as
/// `handled_by`'s do. `handled_by` must still hold a desk of its
/// own - someone has to be answerable for the record.
#[serde(default)]
pub also_moved_by: Vec<String>,
}
fn root() -> String {
"/".to_string()
}
fn default_max_steps() -> usize {
2
}
#[derive(Debug, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct Persona {
pub id: String,
pub need: String,
/// The visitor's situation in their own words - never the page's.
pub text: String,
}
pub fn parse(raw: &str) -> anyhow::Result<NeedsFile> {
serde_yaml::from_str(raw).map_err(|e| anyhow::anyhow!("parsing needs.yaml: {e}"))
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum Level {
Fail,
Warn,
}
#[derive(Debug)]
pub struct Finding {
pub need: String,
pub level: Level,
pub message: String,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct Step {
pub page: String,
pub alternative: String,
}
/// One alternative on a page that still meets the need, with the
/// cheapest way on from it.
struct Viable {
alternative: String,
path: Vec<Step>,
fields: usize,
records_with_action: Option<bool>,
}
type Questions = HashMap<String, Question>;
fn visible(q: &Question, as_group: Option<&str>) -> bool {
match &q.qualifies {
None => true,
Some(group) => as_group == Some(group.as_str()),
}
}
fn required_fields(alt: &Alternative) -> usize {
alt.features
.iter()
.flat_map(|f| &f.requirements)
.filter(|r| !r.optional)
.count()
}
/// The alternatives on `page` from which the need is still met within
/// `remaining` submissions.
fn viable(
questions: &Questions,
need: &Need,
page: &Question,
remaining: usize,
entry: bool,
include_disabled: bool,
trail: &mut Vec<String>,
) -> Vec<Viable> {
let mut out = Vec::new();
if remaining == 0 {
return out;
}
for alt in page.alternatives.iter().filter(|a| include_disabled || !a.disabled) {
if entry && need.via.as_deref().is_some_and(|via| via != alt.name) {
continue;
}
let step = Step {
page: page.id.clone(),
alternative: alt.name.clone(),
};
let fields = required_fields(alt);
if alt.record_as.as_deref() == Some(need.lands_in.as_str()) {
out.push(Viable {
alternative: alt.name.clone(),
path: vec![step],
fields,
records_with_action: Some(alt.action.is_some()),
});
continue;
}
let Some(target) = alt
.action
.as_deref()
.and_then(|a| action_target_id(questions, a))
.and_then(|id| questions.get(&id))
else {
continue;
};
if trail.contains(&target.id) || !visible(target, need.enters_as.as_deref()) {
continue;
}
trail.push(target.id.clone());
let onward = viable(questions, need, target, remaining - 1, false, include_disabled, trail);
trail.pop();
if let Some(best) = onward.into_iter().min_by_key(|v| (v.path.len(), v.fields)) {
let mut path = vec![step];
path.extend(best.path);
out.push(Viable {
alternative: alt.name.clone(),
path,
fields: fields + best.fields,
records_with_action: best.records_with_action,
});
}
}
out
}
fn entry_viable(questions: &Questions, need: &Need, include_disabled: bool) -> Option<Vec<Viable>> {
let page = questions.get(&need.enters)?;
let mut trail = vec![page.id.clone()];
Some(viable(questions, need, page, need.max_steps, true, include_disabled, &mut trail))
}
/// The cheapest path that meets the need, if any.
pub fn path_for(questions: &Questions, need: &Need) -> Option<Vec<Step>> {
entry_viable(questions, need, false)?
.into_iter()
.min_by_key(|v| (v.path.len(), v.fields))
.map(|v| v.path)
}
struct Desk {
page: String,
group: Option<String>,
/// (from, to, button label)
moves: Vec<(String, String, String)>,
}
fn desks_over(questions: &Questions, bucket: &str) -> Vec<Desk> {
let mut desks = Vec::new();
for q in questions.values() {
for resource in q
.alternatives
.iter()
.flat_map(|a| &a.features)
.filter_map(|f| f.resource.as_ref())
.filter(|r| r.bucket() == Some(bucket))
{
desks.push(Desk {
page: q.id.clone(),
group: resource.requires_group.clone().or_else(|| q.qualifies.clone()),
moves: resource
.transitions
.iter()
.map(|t| (t.from.clone(), t.to.clone(), t.label.clone()))
.collect(),
});
}
}
desks.sort_by(|a, b| a.page.cmp(&b.page));
desks
}
/// Self-service moves into `bucket` - a visitor holding their own
/// record's link, no group involved (unsubscribe).
fn self_moves(questions: &Questions, bucket: &str) -> Vec<String> {
questions
.values()
.flat_map(|q| &q.alternatives)
.filter_map(|a| a.self_transition.as_ref())
.filter(|s| s.bucket == bucket)
.map(|s| s.to.clone())
.collect()
}
/// Structural check of every need. `Fail` findings mean the business
/// need is not met by the content as written.
pub fn check(
file: &NeedsFile,
questions: &Questions,
aggregates: &HashMap<String, AggregateSchema>,
) -> Vec<Finding> {
let mut findings = Vec::new();
let mut seen = HashSet::new();
for need in &file.needs {
let mut report = |level: Level, message: String| {
findings.push(Finding {
need: need.id.clone(),
level: if need.planned { Level::Warn } else { level },
message: if need.planned && level == Level::Fail {
format!("planned, not met yet: {message}")
} else {
message
},
})
};
if !seen.insert(need.id.as_str()) {
report(Level::Fail, "declared more than once".to_string());
continue;
}
// The visitor's side: can they get their answer recorded?
match questions.get(&need.enters) {
None => report(Level::Fail, format!("entry page {:?} does not exist", need.enters)),
Some(page) if !visible(page, need.enters_as.as_deref()) => report(
Level::Fail,
format!(
"entry page {:?} is gated on {:?}, which this visitor does not hold",
need.enters, page.qualifies
),
),
Some(page) => {
if page.responsible.is_none() {
report(
Level::Warn,
format!("entry page {:?} names no one responsible to turn to", need.enters),
);
}
let best = entry_viable(questions, need, false)
.unwrap_or_default()
.into_iter()
.min_by_key(|v| (v.path.len(), v.fields));
match best {
None if entry_viable(questions, need, true).is_some_and(|v| !v.is_empty()) => {
let blocked = entry_viable(questions, need, true)
.unwrap_or_default()
.into_iter()
.min_by_key(|v| (v.path.len(), v.fields))
.map(|v| v.path)
.unwrap_or_default();
let at = blocked
.iter()
.find(|s| {
questions[&s.page]
.alternatives
.iter()
.any(|a| a.name == s.alternative && a.disabled)
})
.map(|s| format!("{:?} on {}", s.alternative, s.page))
.unwrap_or_default();
report(
Level::Fail,
format!(
"the path exists but is closed: {at} is `disabled`, so the answer is advertised and refused"
),
)
}
None => report(
Level::Fail,
format!(
"no path from {:?}{} records into {:?} within {} submission(s)",
need.enters,
need.via
.as_ref()
.map(|v| format!(" via {v:?}"))
.unwrap_or_default(),
need.lands_in,
need.max_steps
),
),
Some(best) => {
if let Some(budget) = need.max_fields {
if best.fields > budget {
report(
Level::Fail,
format!(
"the path asks for {} required field(s), over the budget of {budget}",
best.fields
),
);
}
}
if best.records_with_action == Some(false) {
report(
Level::Warn,
"the recording alternative has no `action` - the visitor is never told what happens next".to_string(),
);
}
}
}
}
}
// The business's side: can the handling group finish it?
let Some(schema) = aggregates.get(&need.lands_in) else {
report(
Level::Fail,
format!("aggregates.yaml declares no state machine for {:?}", need.lands_in),
);
continue;
};
let desks = desks_over(questions, &need.lands_in);
let own: Vec<&Desk> = desks
.iter()
.filter(|d| {
d.group.as_deref().is_some_and(|g| g == need.handled_by || need.also_moved_by.iter().any(|m| m == g))
})
.collect();
if !desks.iter().any(|d| d.group.as_deref() == Some(need.handled_by.as_str())) {
let others: Vec<String> = desks
.iter()
.map(|d| format!("{} ({})", d.page, d.group.as_deref().unwrap_or("ungated")))
.collect();
report(
Level::Fail,
format!(
"no desk gated on {:?} reads {:?}{}",
need.handled_by,
need.lands_in,
if others.is_empty() {
String::new()
} else {
format!(" - it is read by {}", others.join(", "))
}
),
);
continue;
}
let selfs = self_moves(questions, &need.lands_in);
let exits = |state: &str| -> Vec<String> {
let allowed = schema.allowed(state);
let mut to: Vec<String> = own
.iter()
.flat_map(|d| &d.moves)
.filter(|(from, to, _)| from == state && allowed.contains(to))
.map(|(_, to, _)| to.clone())
.collect();
to.extend(selfs.iter().filter(|to| allowed.contains(to)).cloned());
to
};
let done: HashSet<&str> = need.done_when.iter().map(String::as_str).collect();
// Reachability walks the whole machine (a finished state may
// lead on to another: subscribed, then unsubscribed). Stranding
// only counts on the way to the first finished state - what
// happens to a handled record afterwards is not this need.
let mut reached: HashSet<String> = HashSet::from([schema.initial.clone()]);
let mut queue = VecDeque::from([(schema.initial.clone(), false)]);
let mut stranded = Vec::new();
while let Some((state, past_done)) = queue.pop_front() {
let past_done = past_done || done.contains(state.as_str());
let next = exits(&state);
if next.is_empty() && !past_done {
stranded.push(state.clone());
}
for to in next {
if reached.insert(to.clone()) {
queue.push_back((to, past_done));
}
}
}
for state in &need.done_when {
if !schema.has_state(state) {
report(
Level::Fail,
format!("done_when names {state:?}, which {:?} has no such state", need.lands_in),
);
} else if !reached.contains(state) {
report(
Level::Fail,
format!(
"{:?} cannot move a record to {state:?} - no desk of theirs offers the moves",
need.handled_by
),
);
}
}
for stage in &need.stages {
if !schema.has_state(stage) {
report(
Level::Fail,
format!(
"the business describes a stage {stage:?} that {:?} has no state for",
need.lands_in
),
);
} else if !reached.contains(stage) {
report(
Level::Fail,
format!(
"stage {stage:?} exists but {:?} is offered no moves that reach it",
need.handled_by
),
);
}
}
for state in &need.success {
if !need.done_when.contains(state) {
report(
Level::Fail,
format!("success names {state:?}, which is not one of done_when"),
);
}
}
stranded.sort();
for state in stranded {
report(
Level::Fail,
format!(
"records strand in {state:?}: {:?} is offered no move out of it",
need.handled_by
),
);
}
}
let need_ids: HashSet<&str> = file.needs.iter().map(|n| n.id.as_str()).collect();
for persona in &file.personas {
if !need_ids.contains(persona.need.as_str()) {
findings.push(Finding {
need: persona.need.clone(),
level: Level::Fail,
message: format!("persona {:?} points at a need that is not declared", persona.id),
});
}
}
findings
}
/// Everything a simulation of the business needs, resolved: each
/// need's path with the choice on every page and the form that gets
/// filled, each bucket's state machine with the buttons every desk
/// offers, and the personas. Portal resolves it because only portal
/// loads content the way the runtime does (derived ids, relative
/// refs, sections); whoever simulates reads plain JSON.
pub fn sim_model(
file: &NeedsFile,
questions: &Questions,
aggregates: &HashMap<String, AggregateSchema>,
) -> serde_json::Value {
use serde_json::json;
let mut needs = Vec::new();
let mut buckets = serde_json::Map::new();
for need in &file.needs {
let Some(path) = path_for(questions, need) else {
continue;
};
let steps: Vec<serde_json::Value> = path
.iter()
.enumerate()
.map(|(i, step)| {
let page = &questions[&step.page];
let mut trail: Vec<String> = path[..=i].iter().map(|s| s.page.clone()).collect();
let expect: Vec<String> =
viable(questions, need, page, need.max_steps - i, i == 0, false, &mut trail)
.into_iter()
.map(|v| v.alternative)
.collect();
let fields: Vec<serde_json::Value> = page
.alternatives
.iter()
.find(|a| a.name == step.alternative)
.into_iter()
.flat_map(|a| &a.features)
.flat_map(|f| &f.requirements)
.map(|r| {
json!({
"name": r.name,
"label": r.label.clone().unwrap_or_else(|| r.name.clone()),
"type": r.kind,
"optional": r.optional,
"options": r.options,
})
})
.collect();
json!({
"page": step.page,
"question": page.name,
"options": page.alternatives.iter().filter(|a| !a.disabled).map(|a| json!({
"name": a.name,
"description": plain(&a.description),
})).collect::<Vec<_>>(),
"expect": expect,
"alternative": step.alternative,
"fields": fields,
})
})
.collect();
needs.push(json!({
"id": need.id,
"who": need.who,
"wants": need.wants,
"enters_as": need.enters_as,
"lands_in": need.lands_in,
"handled_by": need.handled_by,
"done_when": need.done_when,
"success": need.success,
"steps": steps,
}));
if let Some(schema) = aggregates.get(&need.lands_in) {
buckets.entry(need.lands_in.clone()).or_insert_with(|| {
json!({
"initial": schema.initial,
"transitions": schema.transitions,
"desks": desks_over(questions, &need.lands_in).iter().map(|d| json!({
"page": d.page,
"group": d.group,
"moves": d.moves.iter().map(|(from, to, label)| json!({
"from": from, "to": to, "label": label,
})).collect::<Vec<_>>(),
})).collect::<Vec<_>>(),
"self_moves": self_moves(questions, &need.lands_in),
})
});
}
}
json!({
"needs": needs,
"buckets": buckets,
"personas": file.personas.iter().map(|p| json!({
"id": p.id, "need": p.need, "text": plain(&p.text),
})).collect::<Vec<_>>(),
})
}
/// One decision for the answering engine - the shape is laya's
/// `predict(state, questions)` input, so an adapter is a loop and
/// nothing more; any other engine reads the same fields.
#[derive(Debug, Serialize)]
pub struct Task {
pub id: String,
pub persona: String,
pub need: String,
pub page: String,
pub state: BTreeMap<&'static str, String>,
pub questions: BTreeMap<&'static str, ChoiceQuestion>,
/// Every alternative that still meets the need from here. The
/// engine never sees this.
pub expect: Vec<String>,
}
#[derive(Debug, Serialize)]
pub struct ChoiceQuestion {
#[serde(rename = "type")]
pub kind: &'static str,
pub instructions: String,
pub criteria: BTreeMap<String, String>,
}
fn plain(text: &str) -> String {
text.split_whitespace().collect::<Vec<_>>().join(" ")
}
/// Every (persona, page-on-their-path) pair with a real choice on it.
/// `skim` shows the engine only each alternative's heading - what a
/// visitor who never reads the description has to go on.
pub fn tasks(file: &NeedsFile, questions: &Questions, skim: bool) -> Vec<Task> {
let mut out = Vec::new();
for persona in &file.personas {
let Some(need) = file.needs.iter().find(|n| n.id == persona.need) else {
continue;
};
let Some(path) = path_for(questions, need) else {
continue;
};
for (i, step) in path.iter().enumerate() {
let page = &questions[&step.page];
let open: Vec<&Alternative> = page.alternatives.iter().filter(|a| !a.disabled).collect();
if open.len() < 2 {
continue;
}
let mut trail: Vec<String> = path[..=i].iter().map(|s| s.page.clone()).collect();
let expect: Vec<String> = viable(questions, need, page, need.max_steps - i, i == 0, false, &mut trail)
.into_iter()
.map(|v| v.alternative)
.collect();
let criteria = open
.iter()
.map(|a| {
let description = plain(&a.description);
(
a.name.clone(),
if skim || description.is_empty() { a.name.clone() } else { description },
)
})
.collect();
out.push(Task {
id: format!("{}@{}", persona.id, step.page),
persona: persona.id.clone(),
need: need.id.clone(),
page: step.page.clone(),
state: BTreeMap::from([("body", plain(&persona.text))]),
questions: BTreeMap::from([(
"pick",
ChoiceQuestion {
kind: "choice",
instructions: format!(
"A visitor reads the page question {:?}. Which alternative do they choose?",
page.name
),
criteria,
},
)]),
expect,
});
}
}
out
}
#[derive(Debug, Deserialize)]
pub struct Answer {
pub id: String,
pub choice: String,
#[serde(default)]
pub confidence: Option<f64>,
}
#[derive(Debug, Default)]
pub struct Score {
pub total: usize,
pub correct: usize,
/// (task id, picked, expected, confidence)
pub misroutes: Vec<(String, String, Vec<String>, Option<f64>)>,
/// Right, but the engine was under `hesitant_below` - wording that
/// works and is still close to a coin flip.
pub hesitant: Vec<(String, f64)>,
pub unanswered: Vec<String>,
}
impl Score {
pub fn accuracy(&self) -> f64 {
if self.total == 0 {
1.0
} else {
self.correct as f64 / self.total as f64
}
}
}
pub fn score(tasks: &[Task], answers: &[Answer], hesitant_below: f64) -> Score {
let by_id: HashMap<&str, &Answer> = answers.iter().map(|a| (a.id.as_str(), a)).collect();
let mut score = Score {
total: tasks.len(),
..Score::default()
};
for task in tasks {
match by_id.get(task.id.as_str()) {
None => score.unanswered.push(task.id.clone()),
Some(answer) if task.expect.contains(&answer.choice) => {
score.correct += 1;
if let Some(c) = answer.confidence.filter(|c| *c < hesitant_below) {
score.hesitant.push((task.id.clone(), c));
}
}
Some(answer) => score.misroutes.push((
task.id.clone(),
answer.choice.clone(),
task.expect.clone(),
answer.confidence,
)),
}
}
score
}
#[cfg(test)]
mod tests {
use super::*;
use portal::aggregates::parse_aggregates_yaml;
use portal::content::build_questions;
fn site(desk_moves: &str) -> (Questions, HashMap<String, AggregateSchema>) {
let files = vec![
(
"index.yaml".to_string(),
"name: What are you here for?\nresponsible: { name: A, contact: a@b.c }\nalternatives:\n - name: Hire us\n description: Bring a problem.\n action: /thanks\n record_as: projects\n features:\n - name: Brief\n requirements:\n - name: email\n type: email\n - name: brief\n type: textarea\n - name: budget\n optional: true\n - name: Follow along\n description: Notes by mail.\n action: /thanks\n record_as: subscribers\n".to_string(),
),
("thanks.yaml".to_string(), "name: Thanks\nfollowup: true\n".to_string()),
(
"desk/index.yaml".to_string(),
format!("name: Desk\nqualifies: owners\nalternatives:\n - name: Projects\n features:\n - name: \"\"\n resource:\n source: {{ kind: kv, bucket: projects }}\n requires_group: owners\n transitions:\n{desk_moves}"),
),
];
let aggregates = parse_aggregates_yaml(
"aggregates:\n - bucket: projects\n initial: open\n states:\n open: { event: submitted }\n quoted: { event: quoted }\n won: { event: won }\n lost: { event: lost }\n transitions:\n open: [quoted, lost]\n quoted: [won, lost]\n - bucket: subscribers\n attended_by: a mailer\n initial: open\n states:\n open: { event: subscribed }\n transitions: {}\n",
)
.unwrap();
(build_questions(&files).unwrap(), aggregates)
}
const FULL_DESK: &str = " - { to: quoted, label: Quote }\n - { to: lost, label: Lose }\n - { from: quoted, to: won, label: Win }\n - { from: quoted, to: lost, label: Lose }\n";
fn needs(extra: &str) -> NeedsFile {
parse(&format!(
"needs:\n - id: hire\n who: a client\n wants: work done\n lands_in: projects\n handled_by: owners\n done_when: [won, lost]\n{extra}personas:\n - id: founder\n need: hire\n text: I need someone to build our booking tool.\n"
))
.unwrap()
}
fn fails(findings: &[Finding]) -> Vec<&str> {
findings
.iter()
.filter(|f| f.level == Level::Fail)
.map(|f| f.message.as_str())
.collect()
}
#[test]
fn a_met_need_is_silent() {
let (questions, aggregates) = site(FULL_DESK);
let findings = check(&needs(" max_fields: 2\n"), &questions, &aggregates);
assert!(findings.is_empty(), "{findings:?}");
}
#[test]
fn effort_budget_counts_required_fields_only() {
let (questions, aggregates) = site(FULL_DESK);
let findings = check(&needs(" max_fields: 1\n"), &questions, &aggregates);
assert_eq!(fails(&findings), vec!["the path asks for 2 required field(s), over the budget of 1"]);
}
#[test]
fn a_desk_missing_a_move_strands_records_and_cannot_finish() {
// The desk can quote, but never close a quoted job.
let (questions, aggregates) = site(" - { to: quoted, label: Quote }\n - { to: lost, label: Lose }\n");
let findings = check(&needs(""), &questions, &aggregates);
assert_eq!(
fails(&findings),
vec![
"\"owners\" cannot move a record to \"won\" - no desk of theirs offers the moves",
"records strand in \"quoted\": \"owners\" is offered no move out of it",
]
);
}
#[test]
fn a_finished_state_may_lead_on_to_another() {
// quoted counts as handled here, yet won must still be
// reachable past it; lost-after-quoted is not this need's
// business, so nothing strands.
let (questions, aggregates) = site(FULL_DESK);
let mut file = needs("");
file.needs[0].done_when = vec!["quoted".to_string(), "won".to_string(), "lost".to_string()];
let findings = check(&file, &questions, &aggregates);
assert!(findings.is_empty(), "{findings:?}");
}
#[test]
fn a_disabled_alternative_is_named_as_the_reason() {
let (mut questions, aggregates) = site(FULL_DESK);
questions.get_mut("/").unwrap().alternatives[0].disabled = true;
let findings = check(&needs(""), &questions, &aggregates);
assert_eq!(
fails(&findings),
vec!["the path exists but is closed: \"Hire us\" on / is `disabled`, so the answer is advertised and refused"]
);
}
#[test]
fn success_must_be_a_finished_state_and_the_sim_model_carries_the_whole_flow() {
let (questions, aggregates) = site(FULL_DESK);
let mut file = needs(" success: [won]\n");
assert!(check(&file, &questions, &aggregates).is_empty());
let model = sim_model(&file, &questions, &aggregates);
let need = &model["needs"][0];
assert_eq!(need["success"][0], "won");
let step = &need["steps"][0];
assert_eq!(step["question"], "What are you here for?");
assert_eq!(step["expect"][0], "Hire us");
assert_eq!(step["options"].as_array().unwrap().len(), 2);
// The form the visitor fills, required and optional alike.
let fields = step["fields"].as_array().unwrap();
assert_eq!(fields.len(), 3);
assert_eq!(fields[2]["optional"], true);
// The desk's buttons, labelled, per state.
let moves = model["buckets"]["projects"]["desks"][0]["moves"].as_array().unwrap();
assert!(moves.iter().any(|m| m["from"] == "quoted" && m["to"] == "won" && m["label"] == "Win"));
assert_eq!(model["buckets"]["projects"]["initial"], "open");
assert_eq!(model["personas"][0]["id"], "founder");
// A stage the business named must exist as a state.
file.needs[0].stages = vec!["open".into(), "surveyed".into(), "quoted".into()];
assert_eq!(
fails(&check(&file, &questions, &aggregates)),
vec!["the business describes a stage \"surveyed\" that \"projects\" has no state for"]
);
file.needs[0].stages = vec!["open".into(), "quoted".into(), "won".into()];
assert!(check(&file, &questions, &aggregates).is_empty());
file.needs[0].success = vec!["quoted".to_string()];
assert_eq!(
fails(&check(&file, &questions, &aggregates)),
vec!["success names \"quoted\", which is not one of done_when"]
);
}
#[test]
fn a_need_can_be_carried_by_more_than_one_groups_desk() {
// owners may only quote or lose; closing a quoted job is the
// site managers' button, on a desk of their own.
let (mut questions, aggregates) = site(" - { to: quoted, label: Quote }\n - { to: lost, label: Lose }\n");
let more = build_questions(&[(
"site/index.yaml".to_string(),
"name: Site desk\nqualifies: site_managers\nalternatives:\n - name: Projects\n features:\n - name: \"\"\n resource:\n source: { kind: kv, bucket: projects }\n requires_group: site_managers\n transitions:\n - { from: quoted, to: won, label: Win }\n - { from: quoted, to: lost, label: Lose }\n".to_string(),
)])
.unwrap();
questions.extend(more);
let mut file = needs("");
// Alone, owners cannot finish it...
assert_eq!(fails(&check(&file, &questions, &aggregates)).len(), 2);
// ...with the site managers' buttons counted, they can.
file.needs[0].also_moved_by = vec!["site_managers".to_string()];
assert!(check(&file, &questions, &aggregates).is_empty());
// But someone has to be answerable: handled_by needs its own desk.
file.needs[0].handled_by = "nobody".to_string();
assert!(fails(&check(&file, &questions, &aggregates))[0].starts_with("no desk gated on \"nobody\""));
}
#[test]
fn the_wrong_group_holding_the_desk_is_a_failure() {
let (questions, aggregates) = site(FULL_DESK);
let mut file = needs("");
file.needs[0].handled_by = "sales".to_string();
let findings = check(&file, &questions, &aggregates);
assert_eq!(
fails(&findings),
vec!["no desk gated on \"sales\" reads \"projects\" - it is read by /desk (owners)"]
);
}
#[test]
fn via_pins_the_entry_alternative() {
let (questions, aggregates) = site(FULL_DESK);
let mut file = needs("");
file.needs[0].via = Some("Follow along".to_string());
let findings = check(&file, &questions, &aggregates);
assert_eq!(
fails(&findings),
vec!["no path from \"/\" via \"Follow along\" records into \"projects\" within 2 submission(s)"]
);
}
#[test]
fn tasks_carry_the_page_as_a_choice_and_score_separates_misroutes() {
let (questions, _) = site(FULL_DESK);
let file = needs("");
let tasks = tasks(&file, &questions, false);
let skimmed = super::tasks(&file, &questions, true);
assert_eq!(skimmed[0].questions["pick"].criteria["Follow along"], "Follow along");
assert_eq!(tasks.len(), 1);
let task = &tasks[0];
assert_eq!(task.id, "founder@/");
assert_eq!(task.expect, vec!["Hire us"]);
let q = &task.questions["pick"];
assert_eq!(q.kind, "choice");
assert_eq!(q.criteria["Follow along"], "Notes by mail.");
assert!(q.instructions.contains("What are you here for?"));
let right = vec![Answer { id: "founder@/".into(), choice: "Hire us".into(), confidence: Some(0.55) }];
let s = score(&tasks, &right, 0.6);
assert_eq!((s.correct, s.misroutes.len(), s.hesitant.len()), (1, 0, 1));
let wrong = vec![Answer { id: "founder@/".into(), choice: "Follow along".into(), confidence: Some(0.9) }];
let s = score(&tasks, &wrong, 0.6);
assert_eq!((s.correct, s.misroutes.len()), (0, 1));
assert_eq!(s.accuracy(), 0.0);
assert_eq!(score(&tasks, &[], 0.6).unanswered, vec!["founder@/"]);
}
}
+105
View File
@@ -0,0 +1,105 @@
//! Makes a simulation's cases real: every case becomes an aggregate in
//! its bucket, created and moved by the same code the running site
//! uses - `answers::store_answer` for the submission,
//! `aggregates::transition` for each desk move, against a real
//! JetStream with its event log and its refusal of undeclared moves.
//! Then the buckets are read back by replaying each record's events,
//! and reported per need against `done_when` and `success`.
//!
//! The report half runs on its own (`--report-only`) against any
//! JetStream, which is how the same scorecard is taken from a live
//! site's buckets.
//!
//! Point it at a throwaway NATS, never a site's own: it writes.
use std::collections::{BTreeMap, HashMap};
use futures::StreamExt;
use portal::events::store::{ensure_stream, load_events};
use crate::local::{load_aggregates_from_dir, load_from_dir, load_needs_from_dir};
use crate::needs;
use portal::{aggregates, answers, content};
pub async fn run(argv: Vec<String>) -> anyhow::Result<()> {
let value = |flag: &str| argv.iter().position(|a| a == flag).and_then(|i| argv.get(i + 1)).cloned();
let usage = "usage: iris replay --path <questions dir> --nats <url> (--cases <cases.jsonl> | --report-only)";
let dir = value("--path").ok_or_else(|| anyhow::anyhow!(usage))?;
let nats_url = value("--nats").ok_or_else(|| anyhow::anyhow!(usage))?;
let cases_path = value("--cases");
if cases_path.is_none() && !argv.iter().any(|a| a == "--report-only") {
anyhow::bail!(usage);
}
let questions = load_from_dir(&dir)?;
let schemas = content::with_builtin_aggregates(load_aggregates_from_dir(&dir)?);
let file = needs::parse(&load_needs_from_dir(&dir)?.ok_or_else(|| anyhow::anyhow!("the repo has no needs.yaml"))?)?;
let js = async_nats::jetstream::new(async_nats::connect(&nats_url).await?);
// The same call the site makes at boot.
ensure_stream(&js).await?;
if let Some(path) = cases_path {
let raw = std::fs::read_to_string(&path).map_err(|e| anyhow::anyhow!("reading {path}: {e}"))?;
let (mut created, mut moved, mut refused) = (0, 0, 0);
for (n, line) in raw.lines().filter(|l| !l.trim().is_empty()).enumerate() {
let case: serde_json::Value = serde_json::from_str(line)?;
let Some(need) = file.needs.iter().find(|need| case["need"] == need.id.as_str()) else { continue };
let Some(schema) = schemas.get(&need.lands_in) else { continue };
// Only a case that got as far as submitting leaves a record.
let states: Vec<&str> = case["states"].as_array().into_iter().flatten().filter_map(|s| s.as_str()).collect();
if states.is_empty() {
continue;
}
let step = needs::path_for(&questions, need).and_then(|p| p.last().cloned());
let (page, alternative) = step.map(|s| (s.page, s.alternative)).unwrap_or_default();
let id = format!("sim-{}-{}-{n}", case["persona"].as_str().unwrap_or("x"), case["run"]);
let now = chrono::Utc::now().timestamp_millis();
answers::store_answer(&js, &schemas, &need.lands_in, id.clone(), &page, &alternative, &case["record"], now).await?;
created += 1;
for target in &states[1..] {
match aggregates::transition(&js, schema, &id, target, serde_json::json!({"by": "simulation"}), now).await {
Ok(_) => moved += 1,
Err(e) => {
refused += 1;
eprintln!("REFUSED by the engine: {id} -> {target}: {e:?}");
break;
}
}
}
}
println!("{created} aggregate(s) created, {moved} transition(s) accepted, {refused} refused by the engine\n");
}
// Read every bucket back the way the engine knows it: by replaying
// each record's own event log.
let mut by_bucket: HashMap<String, BTreeMap<String, usize>> = HashMap::new();
for bucket in file.needs.iter().map(|n| n.lands_in.clone()).collect::<std::collections::BTreeSet<_>>() {
let (Some(schema), Ok(store)) = (schemas.get(&bucket), js.get_key_value(&bucket).await) else { continue };
let mut keys = store.keys().await?;
while let Some(Ok(id)) = keys.next().await {
let events = load_events(&js, &bucket, &id).await?;
if let Some(record) = aggregates::replay(schema, &id, &events) {
*by_bucket.entry(bucket.clone()).or_default().entry(record.state).or_default() += 1;
}
}
}
println!("{:<18} {:>8} {:>9} {:>9} {}", "bucket", "records", "finished", "success", "by state");
let mut buckets: Vec<_> = by_bucket.iter().collect();
buckets.sort();
for (bucket, states) in buckets {
let needs_here: Vec<&needs::Need> = file.needs.iter().filter(|n| &n.lands_in == bucket).collect();
let is = |pick: fn(&needs::Need) -> &Vec<String>, state: &String| needs_here.iter().any(|n| pick(n).contains(state));
let total: usize = states.values().sum();
let finished: usize = states.iter().filter(|(s, _)| is(|n| &n.done_when, s)).map(|(_, c)| c).sum();
let success: usize = states.iter().filter(|(s, _)| is(|n| &n.success, s)).map(|(_, c)| c).sum();
let pct = |n: usize| if total == 0 { "-".to_string() } else { format!("{:.0}%", n as f64 * 100.0 / total as f64) };
println!(
"{:<18} {:>8} {:>9} {:>9} {}",
bucket,
total,
pct(finished),
pct(success),
states.iter().map(|(s, c)| format!("{s} {c}")).collect::<Vec<_>>().join(", ")
);
}
Ok(())
}