fn parse_task_id(s: &str) -> TaskId {
if s.is_empty() {
return TaskId::named(s);
}
let mut segments = Vec::new();
for part in s.split('.') {
if part.is_empty() {
return TaskId::named(s);
}
if let Ok(n) = part.parse::<u32>() {
segments.push(rhei_core::ast::TaskIdSegment::Number(n));
} else {
segments.push(rhei_core::ast::TaskIdSegment::Named(part.to_string()));
}
}
TaskId::from_segments(segments)
}
fn insert_task_assignee(raw: &str, task_id: &str, assignee: &str) -> MietteResult<String> {
let lines: Vec<&str> = raw.lines().collect();
let mut result: Vec<String> = Vec::with_capacity(lines.len() + 1);
let mut in_target_task = false;
let mut last_metadata_idx: Option<usize> = None;
let mut already_present = false;
let mut inserted = false;
let mut in_code_block = false;
for line in lines.iter() {
if let Some(id) = node_heading_id_outside_code(line, &mut in_code_block) {
if let Some(meta_idx) = last_metadata_idx.take() {
insert_after(&mut result, meta_idx, &format_assignee(assignee));
inserted = true;
}
in_target_task = id == task_id;
}
if !in_code_block && in_target_task && line.starts_with("**Assignee:**") {
already_present = true;
}
if !in_code_block
&& in_target_task
&& (line.starts_with("**State:**") || line.starts_with("**Prior:**"))
{
last_metadata_idx = Some(result.len());
}
result.push((*line).to_string());
}
if already_present {
return Err(miette!(
help = "someone already claimed it. Release it by deleting the **Assignee:** line, \
or claim a different task.",
"Task {} already has an **Assignee:** line",
task_id
));
}
if inserted {
let mut output = result.join("\n");
if raw.ends_with('\n') {
output.push('\n');
}
return Ok(output);
}
let Some(meta_idx) = last_metadata_idx else {
return Err(miette!(
help = "every task needs a `**State:** <state>` line under its heading. Add one, \
then re-run: rhei validate <plan>",
"could not find **State:**/**Prior:** metadata line for Task {} in the markdown",
task_id
));
};
insert_after(&mut result, meta_idx, &format_assignee(assignee));
let mut output = result.join("\n");
if raw.ends_with('\n') {
output.push('\n');
}
Ok(output)
}
fn node_heading_id_outside_code<'a>(
line: &'a str,
in_code_block: &mut bool,
) -> Option<&'a str> {
node_heading_outside_code(line, in_code_block).map(|(_, id)| id)
}
fn node_heading_outside_code<'a>(
line: &'a str,
in_code_block: &mut bool,
) -> Option<(usize, &'a str)> {
if line.trim_start().starts_with("```") {
*in_code_block = !*in_code_block;
return None;
}
if *in_code_block {
return None;
}
node_heading(line)
}
fn node_heading(line: &str) -> Option<(usize, &str)> {
let hashes = line.as_bytes().iter().take_while(|byte| **byte == b'#').count();
if !(3..=6).contains(&hashes) || !line.as_bytes().get(hashes).is_some_and(|b| *b == b' ') {
return None;
}
let body = &line[hashes + 1..];
let (prefix, _) = body.split_once(':')?;
let (_, id) = prefix.rsplit_once(' ')?;
if id.is_empty() { None } else { Some((hashes, id)) }
}
fn format_assignee(value: &str) -> String {
format!("**Assignee:** {}", value)
}
fn insert_after(lines: &mut Vec<String>, idx: usize, value: &str) {
let insert_at = idx + 1;
if insert_at >= lines.len() {
lines.push(value.to_string());
} else {
lines.insert(insert_at, value.to_string());
}
}
#[cfg(test)]
mod next_assignee_rewrite_tests {
use super::*;
#[test]
fn insert_assignee_after_state_when_no_prior() {
let raw = "# Rhei: Test\n\n## Tasks\n\n### Task 1: Work\n**State:** pending\nBody\n";
let rewritten = insert_task_assignee(raw, "1", "codex").expect("rewrite");
assert!(rewritten.contains("**State:** pending\n**Assignee:** codex\nBody"));
}
#[test]
fn insert_assignee_after_prior_when_present() {
let raw =
"# Rhei: Test\n\n## Tasks\n\n### Task 2: Work\n**State:** pending\n**Prior:** Task 1\nBody\n";
let rewritten = insert_task_assignee(raw, "2", "codex").expect("rewrite");
assert!(rewritten.contains("**Prior:** Task 1\n**Assignee:** codex\nBody"));
}
#[test]
fn insert_assignee_supports_child_task_heading() {
let raw = "# Rhei: Test\n\n## Tasks\n\n### Task 1: Parent\n**State:** pending\n\n#### Task 1.1: Child\n**State:** pending\nBody\n";
let rewritten = insert_task_assignee(raw, "1.1", "codex").expect("rewrite");
assert!(rewritten.contains("#### Task 1.1: Child\n**State:** pending\n**Assignee:** codex\nBody"));
assert!(!rewritten.contains("### Task 1: Parent\n**State:** pending\n**Assignee:**"));
}
#[test]
fn insert_assignee_supports_custom_node_kind() {
let raw = "# Rhei: Test\n\n## Tasks\n\n### Bug cache-key: Fix cache\n**State:** pending\nBody\n";
let rewritten = insert_task_assignee(raw, "cache-key", "codex").expect("rewrite");
assert!(rewritten.contains("### Bug cache-key: Fix cache\n**State:** pending\n**Assignee:** codex\nBody"));
}
#[test]
fn insert_assignee_rejects_existing_assignee() {
let raw = "# Rhei: Test\n\n## Tasks\n\n### Task 1: Work\n**State:** pending\n**Assignee:** alice\nBody\n";
let err = insert_task_assignee(raw, "1", "codex").expect_err("existing assignee");
assert!(err.to_string().contains("already has an **Assignee:** line"));
}
#[test]
fn rewrite_state_supports_child_task_heading() {
let raw = "# Rhei: Test\n\n## Tasks\n\n### Task 1: Parent\n**State:** draft\n\n#### Task 1.1: Child\n**State:** draft\nBody\n";
let rewritten = rewrite_task_state(raw, "1.1", "pending").expect("rewrite");
assert!(rewritten.contains("### Task 1: Parent\n**State:** draft"));
assert!(rewritten.contains("#### Task 1.1: Child\n**State:** pending\nBody"));
}
#[test]
fn insert_assignee_ignores_task_shaped_heading_inside_code_fence() {
let raw = "# Rhei: Test\n\n## Tasks\n\n### Task 1: Parent\n**State:** pending\n```markdown\n#### Task 1.1: Example\n**State:** draft\n```\n\n#### Task 1.1: Real child\n**State:** pending\nBody\n";
let rewritten = insert_task_assignee(raw, "1.1", "codex").expect("rewrite");
assert!(rewritten.contains("#### Task 1.1: Example\n**State:** draft\n```"));
assert!(rewritten.contains("#### Task 1.1: Real child\n**State:** pending\n**Assignee:** codex\nBody"));
}
#[test]
fn rewrite_state_ignores_task_shaped_heading_inside_code_fence() {
let raw = "# Rhei: Test\n\n## Tasks\n\n### Task 1: Parent\n**State:** draft\n```markdown\n#### Task 1.1: Example\n**State:** draft\n```\n\n#### Task 1.1: Real child\n**State:** draft\nBody\n";
let rewritten = rewrite_task_state(raw, "1.1", "pending").expect("rewrite");
assert!(rewritten.contains("#### Task 1.1: Example\n**State:** draft\n```"));
assert!(rewritten.contains("#### Task 1.1: Real child\n**State:** pending\nBody"));
}
}
fn rewrite_task_state(raw: &str, task_id: &str, new_state: &str) -> MietteResult<String> {
let lines: Vec<&str> = raw.lines().collect();
let mut result = Vec::with_capacity(lines.len());
let mut in_target_task = false;
let mut state_replaced = false;
let mut in_code_block = false;
for line in &lines {
if !state_replaced {
if let Some(id) = node_heading_id_outside_code(line, &mut in_code_block) {
in_target_task = id == task_id;
}
}
if !in_code_block && in_target_task && !state_replaced && line.starts_with("**State:**") {
let formatted = format!("**State:** {}", format_state_metadata_value(new_state));
result.push(formatted);
state_replaced = true;
continue;
}
result.push(line.to_string());
}
if !state_replaced {
return Err(miette!(
help = "add a `**State:** <state>` line under the task heading, then re-run: \
rhei validate <plan>",
"could not find **State:** line for Task {} in the markdown",
task_id
));
}
let mut output = result.join("\n");
if raw.ends_with('\n') {
output.push('\n');
}
Ok(output)
}
fn next_command(
input: &Path,
state_machine_path: Option<&Path>,
task_id_filter: Option<&str>,
as_json: bool,
no_callbacks: bool,
peek: bool,
rhei_scope: &[String],
) -> MietteResult<()> {
let input_buf = normalize_workspace_input(input);
let input = input_buf.as_path();
let loaded = load_plan(input)?;
let scope = resolve_rhei_scope(&loaded, rhei_scope)?;
let resolved = resolve_state_machines_for_loaded_plan(input, &loaded, state_machine_path)?;
let machines = ExecutionMachines::build(&resolved, input)?;
let workspace_root = execution_workspace_root(&machines.default_callbacks.plan_path);
let report = rhei_validator::validate_with_machine_set(&loaded.rhei, &machines.set);
if report.has_errors() {
return Err(validation_report(input, resolved.default.path.as_deref(), &report.errors));
}
let resolved_filter = task_id_filter
.map(|tid| resolve_cli_task_id(&loaded, tid, &scope))
.transpose()?;
let (task_id_str, current_state_raw, current_state, task_workspace_root) = if let Some(tid) = resolved_filter.as_deref() {
let target_id = parse_task_id(tid);
let task = find_task_by_id(&loaded.rhei.tasks, &target_id)
.ok_or_else(|| {
miette!(
help = format!(
"list the task ids in this plan with: rhei list {}",
shell_quote(&input.display().to_string())
),
"task '{}' not found in the plan",
tid
)
})?;
let open_descendants = open_descendant_tasks(task, &machines.set);
if !open_descendants.is_empty() {
let claimable = narrow_to_rhei_scope(
find_claimable_tasks(
&loaded.rhei,
&machines.set,
&workspace_root,
&loaded.task_roots,
),
&scope,
);
let next_step = match claimable.first() {
Some(candidate) => format!(
"claim what is ready instead: rhei next {} --task {}",
shell_quote(&input.display().to_string()),
candidate.id
),
None => format!(
"finish or cancel the open descendants first, then claim this ticket. \
See every task and its state with: rhei list {}",
shell_quote(&input.display().to_string())
),
};
return Err(miette!(
help = next_step,
"Task {} cannot be claimed while {} descendant task(s) are still open.\n\
Open descendants: {}",
tid,
open_descendants.len(),
format_open_descendants(&open_descendants, &machines.set)
));
}
if let Some(assignee) = task.assignee.as_deref() {
return Err(miette!(
help = format!(
"release it by deleting the **Assignee:** line from Task {tid}, or claim \
whatever is ready instead: rhei next {}",
shell_quote(&input.display().to_string())
),
"Task {} is already assigned to {}",
tid,
assignee
));
}
let machine = machines.for_task_str(tid);
let state_name = normalized_state_name(task.state.as_str(), machine);
let is_initial = task_is_in_initial_state(task, &state_name, machine);
if is_initial {
let mut all_tasks = Vec::new();
collect_plan_tasks(&loaded.rhei.tasks, &mut all_tasks);
let state_map = plan_state_map(&all_tasks, &machines.set);
let all_priors_done = task.prior.iter().all(|dep_id| {
state_map
.get(dep_id)
.map(|s| dependency_is_satisfied(s, machines.set.for_task(dep_id)))
.unwrap_or(false)
});
if !all_priors_done {
let detail = first_blocking_prior(task, &state_map, &machines.set, &scope)
.map(|prior| format!("; waiting on {}", prior))
.unwrap_or_default();
return Err(miette!(
help = format!(
"finish the prerequisite first, or see what is claimable now: rhei list {}",
shell_quote(&input.display().to_string())
),
"Task {} is blocked by incomplete prerequisites{}",
tid,
detail
));
}
}
let state_def = machine
.states
.get(&state_name)
.ok_or_else(|| {
miette!(help = internal_error_help(), "state '{}' missing from loaded machine", state_name)
})?;
let settings = load_merged_settings(&workspace_root)?;
let task_workspace_root = loaded.task_root(tid, &workspace_root);
ensure_state_inputs_exist_for_transition(
&task_workspace_root,
Some(task),
tid,
&state_name,
state_def,
Some(render_visit_count(
loaded.rhei.metadata.as_ref(),
&task.id,
&state_name,
task.state.as_str(),
machine,
)),
machine,
&settings,
&format!("Task {} cannot be claimed in state {}.", tid, state_name),
)?;
(tid.to_string(), task.state.as_str().to_string(), state_name, task_workspace_root)
} else {
let ready = narrow_to_rhei_scope(
find_claimable_tasks(&loaded.rhei, &machines.set, &workspace_root, &loaded.task_roots),
&scope,
);
if ready.is_empty() {
return Err(miette!(
help = "see every task and its state with: rhei list <plan>",
"{}",
diagnose_no_claimable(
&loaded.rhei,
&machines.set,
input,
resolved.default.path.as_deref(),
&scope
)
));
}
let task = ready.into_iter().next().unwrap();
let machine = machines.for_task(&task.id);
let state_name = normalized_state_name(task.state.as_str(), machine);
let state_def = machine
.states
.get(&state_name)
.ok_or_else(|| {
miette!(help = internal_error_help(), "state '{}' missing from loaded machine", state_name)
})?;
let settings = load_merged_settings(&workspace_root)?;
let task_workspace_root = loaded.task_root(&task.id.to_string(), &workspace_root);
ensure_state_inputs_exist_for_transition(
&task_workspace_root,
Some(task),
&task.id.to_string(),
&state_name,
state_def,
Some(render_visit_count(
loaded.rhei.metadata.as_ref(),
&task.id,
&state_name,
task.state.as_str(),
machine,
)),
machine,
&settings,
&format!("Task {} cannot be claimed in state {}.", task.id, state_name),
)?;
(task.id.to_string(), task.state.to_string(), state_name, task_workspace_root)
};
let target_id = parse_task_id(&task_id_str);
let machine = machines.for_task_str(&task_id_str);
let callback_paths = machines.callbacks_for_str(&task_id_str);
let selected_task = find_task_by_id(&loaded.rhei.tasks, &target_id)
.ok_or_else(|| {
miette!(
help = format!(
"list the task ids in this plan with: rhei list {}",
shell_quote(&input.display().to_string())
),
"task '{}' not found in the plan",
task_id_str
)
})?;
let is_initial = task_is_in_initial_state(selected_task, ¤t_state, machine);
let current_state_def = machine
.states
.get(¤t_state)
.ok_or_else(|| {
miette!(help = internal_error_help(), "state '{}' missing from loaded machine", current_state)
})?;
let auto_transition_initial = is_initial
&& !state_declares_autonomous_execution(current_state_def)
&& initial_state_has_non_terminal_forward_transition(selected_task, &loaded.rhei, machine)?;
let route = loaded.task_route(&task_id_str, input);
let final_state = if auto_transition_initial && !peek {
let target_id = parse_task_id(&task_id_str);
let task = find_task_by_id(&loaded.rhei.tasks, &target_id)
.ok_or_else(|| {
miette!(
help = format!(
"list the task ids in this plan with: rhei list {}",
shell_quote(&input.display().to_string())
),
"task '{}' not found in the plan",
task_id_str
)
})?;
let to_state = find_next_transition(task, &loaded.rhei, machine)?.ok_or_else(|| {
miette!(
help = format!(
"no transition leaves '{current_state_raw}'. See the machine's edges with: \
rhei states"
),
"no forward transition available from state '{}'",
current_state_raw
)
})?;
execute_transition(
TransitionFiles { task_file: &route.task_file, metadata_file: &route.metadata_file, metadata_id: &route.metadata_id, artifact_root: &route.execution_root, artifact_id: &task_id_str },
callback_paths,
machine,
&route.local_id,
¤t_state,
&to_state,
None,
no_callbacks,
)?
} else {
current_state.clone()
};
let loaded = load_plan(input)?;
let target_id = parse_task_id(&task_id_str);
let task = find_task_by_id(&loaded.rhei.tasks, &target_id)
.ok_or_else(|| {
miette!(help = internal_error_help(), "task '{}' not found after transition", task_id_str)
})?;
let settings = load_merged_settings(&workspace_root)?;
let no_agent_opts = default_run_options();
let resolved = match resolve_agent_for_task(machine, &final_state, &settings, &no_agent_opts, task) {
Ok(resolved) => resolved,
Err(err) => {
eprintln!(
"warning: could not resolve agent for state '{}': {}",
final_state, err
);
None
}
};
let agent_id_str = resolved.as_ref().map(|r| r.agent.id().to_string());
let model_id_str = resolved.as_ref().and_then(|r| r.model.clone());
let model_provider_str = resolved.as_ref().and_then(|r| r.model_provider.clone());
let model_name_str = resolved.as_ref().and_then(|r| r.model_name.clone());
let mut claimed_as: Option<String> = None;
if !peek && task.assignee.is_none() {
let assignee = agent_id_str.as_deref().unwrap_or("manual");
claimed_as = Some(assignee.to_string());
let final_state_def = machine
.states
.get(&final_state)
.ok_or_else(|| miette!(
help = internal_error_help(),
"state '{}' missing from loaded machine", final_state
))?;
write_task_assignee(
&route.task_file,
&route.local_id,
&task_id_str,
&final_state,
machine,
TaskAssigneeClaimContext {
workspace_root: &task_workspace_root,
metadata: loaded.rhei.metadata.as_ref(),
state_def: final_state_def,
settings: &settings,
},
assignee,
)?;
}
let tooling = resolve_tooling(machine, &final_state, &settings);
let render_context = RuntimeTemplateContext {
workspace_root: &task_workspace_root,
task_roots: Some(&loaded.task_roots),
checkout_root: &task_workspace_root,
plan_path: &callback_paths.plan_path,
state_machine_path: callback_paths.state_machine_path.as_deref(),
plan_title: &loaded.rhei.title,
task,
state_name: &final_state,
current_state_raw: task.state.as_str(),
machine,
metadata: loaded.rhei.metadata.as_ref(),
target: resolved.as_ref().and_then(|r| r.target.as_ref()),
model: model_id_str.as_deref(),
model_provider: model_provider_str.as_deref(),
model_name: model_name_str.as_deref(),
agent: agent_id_str.as_deref(),
agent_mode: resolved.as_ref().and_then(|r| r.mode.as_deref()),
tooling: Some(&tooling),
};
let instructions = resolve_runtime_template_text(
state_instructions(machine, &final_state).as_str(),
&render_context,
);
let personality = state_personality(machine, final_state.as_str())
.map(|text| resolve_runtime_template_text(&text, &render_context));
print_next_output(NextOutput {
as_json,
peek,
claimed_as: claimed_as.as_deref(),
task,
from_state: ¤t_state_raw,
to_state: task.state.as_str(),
personality: personality.as_deref(),
instructions: &instructions,
agent_id: agent_id_str.as_deref(),
model_id: model_id_str.as_deref(),
});
Ok(())
}