use std::path::Path;
use oneagentgraph::config::ConfigRef;
use onevcs::registry::{RepoType, Workflow};
use onevcs::MergePolicy;
use serde::{Deserialize, Serialize};
use crate::error::{Error, Result};
pub const PLAN_SCHEMA_VERSION: u32 = 3;
pub const PLAN_SCHEMA_VERSIONS_READ: [u32; 3] = [PLAN_SCHEMA_VERSION, 2, 1];
pub(crate) const TITLE_IS_REQUIRED: &str = "a lifecycle node states the title its change request \
opens under, and this one names no `title`";
pub(crate) fn body_is_newer(declared: u32) -> String {
format!(
"`body` is a schema {PLAN_SCHEMA_VERSION} field and this plan declares schema_version \
{declared} — set `schema_version: {PLAN_SCHEMA_VERSION}`"
)
}
pub const PLANNER_CONTEXT_HEADING: &str = "## Planner context";
pub(crate) const DONE_WHEN_RETIRED: &str =
"`done_when` is no longer a plan field. A node's review bar is the \
`## Acceptance criteria` section of its own task, which the judge is handed \
verbatim; a bar broader than one node belongs in the onejudge base config the \
node-scope graph's worker already points at, under `user.done_when`";
const DONE_WHEN: &str = "done_when";
pub(crate) const NO_GOAL: &str = "(no goal stated)";
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct Plan {
pub schema_version: u32,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub goal: Option<Goal>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub name: Option<String>,
#[serde(default = "default_concurrency")]
pub concurrency: u32,
pub tasks: Vec<Node>,
}
fn default_concurrency() -> u32 {
4
}
impl Plan {
pub fn load(path: &Path) -> Result<Self> {
let text = std::fs::read_to_string(path).map_err(|e| Error::Ledger {
path: path.to_path_buf(),
source: e,
})?;
let is_json = path
.extension()
.is_some_and(|ext| ext.eq_ignore_ascii_case("json"));
let named = |e: String| Error::Invalid(format!("{}: {e}", path.display()));
let refused = |e: String| named(retired_field_refusal_in(&text).unwrap_or(e));
if is_json {
match serde_json::from_str::<serde_json::Value>(&text) {
Ok(serde_json::Value::Object(_)) => {
return serde_json::from_str(&text).map_err(|e| refused(e.to_string()));
}
Ok(other) => {
let kind = match other {
serde_json::Value::Array(_) => "list",
serde_json::Value::String(_) => "string",
serde_json::Value::Null => "null",
_ => "scalar",
};
return Err(named(format!("must be a JSON mapping, got {kind}")));
}
Err(_) => {}
}
}
serde_norway::from_str(&text).map_err(|e| refused(e.to_string()))
}
}
pub(crate) fn retired_field_refusal(document: &serde_json::Value) -> Option<String> {
match document {
serde_json::Value::Object(map) => {
if map.contains_key(DONE_WHEN) {
let whose = map
.get("id")
.and_then(serde_json::Value::as_str)
.map(|id| format!("'{id}': "))
.unwrap_or_default();
return Some(format!("{whose}{DONE_WHEN_RETIRED}"));
}
map.values().find_map(retired_field_refusal)
}
serde_json::Value::Array(items) => items.iter().find_map(retired_field_refusal),
_ => None,
}
}
fn retired_field_refusal_in(text: &str) -> Option<String> {
retired_field_refusal(&serde_norway::from_str::<serde_json::Value>(text).ok()?)
}
impl Node {
pub fn rendered_task(&self) -> String {
render_task(
self.task.as_deref().unwrap_or_default(),
self.context.as_deref(),
)
}
}
impl Step {
pub fn rendered_task(&self, node_context: Option<&str>) -> String {
render_task(self.task.as_deref().unwrap_or_default(), node_context)
}
}
fn render_task(task: &str, context: Option<&str>) -> String {
match context.map(str::trim).filter(|note| !note.is_empty()) {
None => task.to_string(),
Some(note) => format!(
"{}\n\n{PLANNER_CONTEXT_HEADING}\n\
This reports observed state and adds no acceptance criteria.\n\n{note}\n",
task.trim_end()
),
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct Goal {
pub text: String,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum NodeKind {
#[default]
Agent,
Human,
}
impl NodeKind {
fn is_agent(&self) -> bool {
matches!(self, Self::Agent)
}
}
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct Node {
pub id: String,
#[serde(default, skip_serializing_if = "NodeKind::is_agent")]
pub kind: NodeKind,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub task: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub persona: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub deps: Vec<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_turns: Option<u32>,
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
pub expects_no_diff: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub context: Option<String>,
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
pub parked: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub executor: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub agent_graph: Option<ConfigRef>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub repo: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub repo_type: Option<RepoType>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub workflow: Option<Workflow>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub merge_policy: Option<MergePolicy>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub base_branch: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub branch: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub title: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub body: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub execution_checkout: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub verify_via_ci: Option<bool>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub steps: Option<Vec<Step>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub resume: Option<Resume>,
}
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct Step {
pub id: String,
#[serde(default, skip_serializing_if = "NodeKind::is_agent")]
pub kind: NodeKind,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub task: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub persona: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub deps: Vec<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_turns: Option<u32>,
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
pub expects_no_diff: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub executor: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub agent_graph: Option<ConfigRef>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct Resume {
pub branch: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub checkpoint: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub completed_steps: Vec<String>,
}
#[cfg(test)]
mod tests {
use super::*;
fn scratch(name: &str) -> std::path::PathBuf {
let dir =
std::env::temp_dir().join(format!("onepipeline-plan-{name}-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).expect("a scratch root");
dir
}
#[test]
fn a_json_plan_is_read_with_json_escape_semantics() {
let root = scratch("json");
let path = root.join("emoji.plan.json");
std::fs::write(
&path,
r#"{"schema_version":2,"tasks":[{"id":"a","persona":"engineer","task":"😀 ship it"}]}"#,
)
.expect("written");
let plan = Plan::load(&path).expect("a JSON plan loads");
assert!(
plan.tasks[0]
.task
.as_deref()
.expect("task")
.starts_with('\u{1f600}'),
"the surrogate pair did not survive as one character"
);
assert_eq!(
plan.concurrency, 4,
"the default concurrency was not applied"
);
std::fs::remove_dir_all(&root).ok();
}
#[test]
fn a_json_document_that_is_not_a_mapping_is_refused_by_name() {
let root = scratch("notmapping");
let path = root.join("list.plan.json");
std::fs::write(&path, "[1, 2, 3]").expect("written");
let message = Plan::load(&path).unwrap_err().to_string();
assert!(
message.contains("must be a JSON mapping, got list"),
"{message}"
);
std::fs::remove_dir_all(&root).ok();
}
#[test]
fn a_json_named_file_that_json_cannot_parse_falls_back_to_yaml() {
let root = scratch("yamlfallback");
let path = root.join("actually.plan.json");
std::fs::write(
&path,
"schema_version: 2\ntasks:\n - id: a\n persona: engineer\n task: do it\n",
)
.expect("written");
let plan = Plan::load(&path).expect("the YAML reading is the fallback");
assert_eq!(plan.tasks[0].id, "a");
std::fs::remove_dir_all(&root).ok();
}
#[test]
fn a_plan_with_an_unknown_field_is_refused_at_its_trust_boundary() {
let root = scratch("unknown");
let path = root.join("typo.plan.json");
std::fs::write(
&path,
r#"{"schema_version":2,"concurency":2,"tasks":[{"id":"a","persona":"e","task":"t"}]}"#,
)
.expect("written");
let message = Plan::load(&path).unwrap_err().to_string();
assert!(message.contains("concurency"), "{message}");
std::fs::remove_dir_all(&root).ok();
}
#[test]
fn a_missing_plan_file_names_the_path_it_could_not_read() {
let message = Plan::load(std::path::Path::new("no/such/plan.json"))
.unwrap_err()
.to_string();
assert!(message.contains("no/such/plan.json"), "{message}");
}
#[test]
fn both_shipped_examples_load_and_keep_what_they_declare() {
let root = std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("examples");
let single = Plan::load(&root.join("single-node.plan.json")).expect("single-node loads");
assert_eq!(single.tasks.len(), 1);
assert!(single.goal.is_some());
let mixed = Plan::load(&root.join("mixed-graph.plan.json")).expect("mixed-graph loads");
assert_eq!(mixed.concurrency, 3);
let docs = mixed
.tasks
.iter()
.find(|n| n.id == "docs")
.expect("the docs node");
assert_eq!(
docs.agent_graph.as_ref().map(|r| r.0.as_str()),
Some("./graphs/node-scope.yaml"),
"the example does not reference the shipped node-scope config"
);
let service = mixed
.tasks
.iter()
.find(|n| n.id == "service")
.expect("the service node");
assert_eq!(service.executor.as_deref(), Some("local"));
assert_eq!(service.steps.as_ref().map(Vec::len), Some(2));
}
#[test]
fn a_planner_note_renders_as_its_own_section_and_disclaims_itself() {
let node = Node {
id: "build".into(),
persona: Some("engineer".into()),
task: Some("## What\nship it".into()),
context: Some("the fixture moved to tests/data".into()),
..Node::default()
};
let rendered = node.rendered_task();
assert!(rendered.starts_with("## What\nship it"), "{rendered}");
assert!(rendered.contains(PLANNER_CONTEXT_HEADING), "{rendered}");
assert!(
rendered.contains("adds no acceptance criteria"),
"{rendered}"
);
assert!(
rendered.contains("the fixture moved to tests/data"),
"{rendered}"
);
}
#[test]
fn a_node_with_no_note_renders_its_task_unchanged() {
let node = Node {
id: "build".into(),
task: Some("## What\nship it".into()),
..Node::default()
};
assert_eq!(node.rendered_task(), "## What\nship it");
let blank = Node {
context: Some(" ".into()),
..node
};
assert_eq!(blank.rendered_task(), "## What\nship it");
}
#[test]
fn a_workstreams_note_reaches_every_agent_step() {
let step = Step {
id: "implement".into(),
persona: Some("engineer".into()),
task: Some("## What\nimplement".into()),
..Step::default()
};
let rendered = step.rendered_task(Some("the API moved"));
assert!(rendered.contains(PLANNER_CONTEXT_HEADING), "{rendered}");
assert_eq!(step.rendered_task(None), "## What\nimplement");
}
#[test]
fn a_plan_round_trips_without_growing_the_fields_it_omitted() {
let source = r#"{"schema_version":2,"tasks":[{"id":"a","persona":"e","task":"t"}]}"#;
let plan: Plan = serde_json::from_str(source).expect("it parses");
let written = serde_json::to_string(&plan).expect("it serialises");
assert!(
!written.contains("\"kind\""),
"{written} grew a default kind"
);
assert!(
!written.contains("\"deps\""),
"{written} grew an empty deps"
);
assert!(
!written.contains("\"parked\""),
"{written} grew a false parked"
);
}
#[test]
fn a_current_version_plan_round_trips_and_omits_the_budget_it_does_not_declare() {
let source = format!(
r#"{{"schema_version":{PLAN_SCHEMA_VERSION},"name":"round-trip","tasks":[
{{"id":"budgeted","persona":"e","task":"t","max_turns":45}},
{{"id":"plain","persona":"e","task":"t"}}]}}"#
);
let plan: Plan = serde_json::from_str(&source).expect("it parses");
assert_eq!(plan.schema_version, PLAN_SCHEMA_VERSION);
assert_eq!(plan.tasks[0].max_turns, Some(45));
assert_eq!(plan.tasks[1].max_turns, None);
let written = serde_json::to_string(&plan).expect("it serialises");
assert!(
written.contains(&format!("\"schema_version\":{PLAN_SCHEMA_VERSION}")),
"the version a reader decides by is not on the wire: {written}"
);
assert_eq!(
written.matches("\"max_turns\"").count(),
1,
"a node that declared no turn budget was written one: {written}"
);
assert!(written.contains("\"max_turns\":45"), "{written}");
assert_eq!(
serde_json::from_str::<Plan>(&written).expect("it re-parses"),
plan,
"the plan did not survive a round trip through this crate"
);
}
#[test]
fn a_plan_carrying_done_when_is_answered_about_the_field_at_every_version() {
let root = scratch("donewhen");
for version in PLAN_SCHEMA_VERSIONS_READ {
let path = root.join(format!("v{version}.plan.json"));
std::fs::write(
&path,
format!(
r#"{{"schema_version":{version},"tasks":[
{{"id":"contract","persona":"e","task":"t",
"done_when":"the gate is green"}}]}}"#
),
)
.expect("written");
let message = Plan::load(&path).unwrap_err().to_string();
assert!(message.contains("'contract':"), "{message}");
assert!(message.contains(DONE_WHEN_RETIRED), "{message}");
assert!(
!message.contains("schema_version"),
"a version refusal displaced the field's: {message}"
);
}
std::fs::remove_dir_all(&root).ok();
}
}