use std::collections::HashMap;
use crate::error::{OrchestratorError, Result};
use crate::hooks::HooksConfig;
use crate::vcs::VcsBackend;
use serde::{Deserialize, Serialize};
use super::defaults::{self, *};
use super::expand;
fn default_suppress_repetitive_debug() -> bool {
DEFAULT_SUPPRESS_REPETITIVE_DEBUG
}
fn default_log_summary_interval_secs() -> u64 {
DEFAULT_LOG_SUMMARY_INTERVAL_SECS
}
fn default_stall_detection_enabled() -> bool {
DEFAULT_STALL_DETECTION_ENABLED
}
fn default_stall_detection_threshold() -> u32 {
DEFAULT_STALL_DETECTION_THRESHOLD
}
fn default_error_circuit_breaker_enabled() -> bool {
DEFAULT_ERROR_CIRCUIT_BREAKER_ENABLED
}
fn default_error_circuit_breaker_threshold() -> usize {
DEFAULT_ERROR_CIRCUIT_BREAKER_THRESHOLD
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct LoggingConfig {
#[serde(default = "default_suppress_repetitive_debug")]
pub suppress_repetitive_debug: bool,
#[serde(default = "default_log_summary_interval_secs")]
pub summary_interval_secs: u64,
}
impl Default for LoggingConfig {
fn default() -> Self {
Self {
suppress_repetitive_debug: DEFAULT_SUPPRESS_REPETITIVE_DEBUG,
summary_interval_secs: DEFAULT_LOG_SUMMARY_INTERVAL_SECS,
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct StallDetectionConfig {
#[serde(default = "default_stall_detection_enabled")]
pub enabled: bool,
#[serde(default = "default_stall_detection_threshold")]
pub threshold: u32,
#[serde(default)]
pub apply_escalation_after_empty_wip: Option<u32>,
#[serde(default)]
pub apply_escalation_max_uses_per_stall: Option<u32>,
}
impl Default for StallDetectionConfig {
fn default() -> Self {
Self {
enabled: DEFAULT_STALL_DETECTION_ENABLED,
threshold: DEFAULT_STALL_DETECTION_THRESHOLD,
apply_escalation_after_empty_wip: DEFAULT_APPLY_ESCALATION_AFTER_EMPTY_WIP,
apply_escalation_max_uses_per_stall: DEFAULT_APPLY_ESCALATION_MAX_USES_PER_STALL,
}
}
}
impl StallDetectionConfig {
pub fn validate(&self) -> Result<()> {
if let Some(after) = self.apply_escalation_after_empty_wip {
if after >= self.threshold {
return Err(OrchestratorError::ConfigLoad(format!(
"stall_detection.apply_escalation_after_empty_wip ({after}) must be less than stall_detection.threshold ({})",
self.threshold
)));
}
}
if matches!(self.apply_escalation_max_uses_per_stall, Some(0)) {
return Err(OrchestratorError::ConfigLoad(
"stall_detection.apply_escalation_max_uses_per_stall must be at least 1 when set"
.to_string(),
));
}
Ok(())
}
pub fn apply_escalation_policy_enabled(&self) -> bool {
self.apply_escalation_after_empty_wip.is_some()
&& self.apply_escalation_max_uses_per_stall.is_some()
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)]
pub struct AcceptanceEscalationConfig {
#[serde(default)]
pub after_invalid_results: Option<u32>,
#[serde(default)]
pub max_uses_per_sequence: Option<u32>,
}
impl AcceptanceEscalationConfig {
pub fn after_invalid_results(&self) -> u32 {
self.after_invalid_results
.unwrap_or(DEFAULT_ACCEPTANCE_ESCALATION_AFTER_INVALID_RESULTS)
}
pub fn max_uses_per_sequence(&self) -> u32 {
self.max_uses_per_sequence
.unwrap_or(DEFAULT_ACCEPTANCE_ESCALATION_MAX_USES_PER_SEQUENCE)
}
fn merge(&mut self, other: Self) {
overwrite_if_some(&mut self.after_invalid_results, other.after_invalid_results);
overwrite_if_some(&mut self.max_uses_per_sequence, other.max_uses_per_sequence);
}
pub fn validate(&self) -> Result<()> {
for (field, value) in [
("after_invalid_results", self.after_invalid_results),
("max_uses_per_sequence", self.max_uses_per_sequence),
] {
if matches!(value, Some(0)) {
return Err(OrchestratorError::ConfigLoad(format!(
"Configuration error: `acceptance_escalation.{field}` must be a positive \
integer (at least 1); remove `acceptance_escalation_command` to disable \
Acceptance escalation instead of setting 0"
)));
}
}
Ok(())
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct ErrorCircuitBreakerConfig {
#[serde(default = "default_error_circuit_breaker_enabled")]
pub enabled: bool,
#[serde(default = "default_error_circuit_breaker_threshold")]
pub threshold: usize,
}
impl Default for ErrorCircuitBreakerConfig {
fn default() -> Self {
Self {
enabled: default_error_circuit_breaker_enabled(),
threshold: default_error_circuit_breaker_threshold(),
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct OrchestratorConfig {
#[serde(default)]
pub apply_command: Option<String>,
#[serde(default)]
pub apply_escalation_command: Option<String>,
#[serde(default)]
pub apply_stall_diagnose_command: Option<String>,
#[serde(default)]
pub archive_command: Option<String>,
#[serde(default)]
pub apply_skill: Option<String>,
#[serde(default)]
pub archive_skill: Option<String>,
#[serde(default)]
pub analyze_skill: Option<String>,
#[serde(default)]
pub accept_skill: Option<String>,
#[serde(default)]
pub rejecting_skill: Option<String>,
#[serde(default)]
pub cleanup_review_skill: Option<String>,
#[serde(default)]
pub resolve_skill: Option<String>,
#[serde(default)]
pub analyze_command: Option<String>,
#[serde(default)]
pub acceptance_command: Option<String>,
#[serde(default)]
pub acceptance_escalation_command: Option<String>,
#[serde(default)]
pub acceptance_escalation: Option<AcceptanceEscalationConfig>,
#[serde(default)]
pub apply_prompt: Option<String>,
#[serde(default)]
pub apply_append_prompt: Option<String>,
#[serde(default)]
pub acceptance_prompt: Option<String>,
#[serde(default)]
pub acceptance_append_prompt: Option<String>,
#[serde(default)]
pub acceptance_prompt_mode: Option<AcceptancePromptMode>,
#[serde(default)]
pub archive_prompt: Option<String>,
#[serde(default)]
pub archive_append_prompt: Option<String>,
#[serde(default)]
pub analyze_append_prompt: Option<String>,
#[serde(default)]
pub resolve_append_prompt: Option<String>,
#[serde(default)]
pub envs: Option<HashMap<String, String>>,
#[serde(default)]
pub hooks: Option<HooksConfig>,
#[serde(default)]
pub logging: Option<LoggingConfig>,
#[serde(default)]
pub stall_detection: Option<StallDetectionConfig>,
#[serde(default)]
pub error_circuit_breaker: Option<ErrorCircuitBreakerConfig>,
#[serde(default)]
pub completion_check_delay_ms: Option<u64>,
#[serde(default)]
pub completion_check_max_retries: Option<u32>,
#[serde(default)]
pub max_iterations: Option<u32>,
#[serde(default)]
pub max_concurrent_workspaces: Option<usize>,
#[serde(default)]
pub workspace_base_dir: Option<String>,
#[serde(default)]
pub state_base_dir: Option<String>,
#[serde(default)]
pub resolve_command: Option<String>,
#[serde(default)]
pub use_llm_analysis: Option<bool>,
#[serde(default)]
pub vcs_backend: Option<VcsBackend>,
#[serde(default)]
pub propose_command: Option<String>,
#[serde(default)]
pub worktree_command: Option<String>,
#[serde(default)]
pub command_queue_stagger_delay_ms: Option<u64>,
#[serde(default)]
pub command_queue_max_retries: Option<u32>,
#[serde(default)]
pub command_queue_retry_delay_ms: Option<u64>,
#[serde(default)]
pub command_queue_retry_patterns: Option<Vec<String>>,
#[serde(default)]
pub command_queue_retry_if_duration_under_secs: Option<u64>,
#[serde(default)]
pub acceptance_max_continues: Option<u32>,
#[serde(default)]
pub command_inactivity_timeout_secs: Option<u64>,
#[serde(default)]
pub command_inactivity_kill_grace_secs: Option<u64>,
#[serde(default)]
pub command_inactivity_timeout_max_retries: Option<u32>,
#[serde(default)]
pub command_max_runtime_secs: Option<u64>,
#[serde(default)]
pub acceptance_max_runtime_secs: Option<u64>,
#[serde(default)]
pub stream_json_textify: Option<bool>,
#[serde(default)]
pub command_strict_process_cleanup: Option<bool>,
#[serde(default)]
pub lifecycle_integration: Option<LifecycleIntegrationConfig>,
#[serde(default)]
pub judge_commands: Option<JudgeCommandsConfig>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)]
pub struct LifecycleIntegrationConfig {
#[serde(default)]
pub enabled: Option<bool>,
#[serde(default)]
pub command: Vec<String>,
#[serde(default)]
pub queue_capacity: Option<usize>,
#[serde(default)]
pub write_timeout_ms: Option<u64>,
#[serde(default)]
pub shutdown_timeout_ms: Option<u64>,
}
impl LifecycleIntegrationConfig {
pub fn is_enabled(&self) -> bool {
match self.enabled {
Some(false) => false,
Some(true) => true,
None => !self.command.is_empty(),
}
}
pub fn queue_capacity(&self) -> usize {
self.queue_capacity
.unwrap_or(DEFAULT_LIFECYCLE_QUEUE_CAPACITY)
}
pub fn write_timeout_ms(&self) -> u64 {
self.write_timeout_ms
.unwrap_or(DEFAULT_LIFECYCLE_WRITE_TIMEOUT_MS)
}
pub fn shutdown_timeout_ms(&self) -> u64 {
self.shutdown_timeout_ms
.unwrap_or(DEFAULT_LIFECYCLE_SHUTDOWN_TIMEOUT_MS)
}
pub fn validate(&self) -> Result<()> {
if !self.is_enabled() {
return Ok(());
}
if self.command.is_empty() {
return Err(OrchestratorError::ConfigLoad(
"Configuration error: `lifecycle_integration.command` must be a non-empty argv array, for example [\"cflx-herdr-adapter\"]".to_string(),
));
}
if self.command[0].trim().is_empty() {
return Err(OrchestratorError::ConfigLoad(
"Configuration error: `lifecycle_integration.command[0]` must be a non-empty executable name".to_string(),
));
}
if matches!(self.queue_capacity, Some(0)) {
return Err(OrchestratorError::ConfigLoad(
"Configuration error: `lifecycle_integration.queue_capacity` must be at least 1"
.to_string(),
));
}
if matches!(self.write_timeout_ms, Some(0)) {
return Err(OrchestratorError::ConfigLoad(
"Configuration error: `lifecycle_integration.write_timeout_ms` must be at least 1"
.to_string(),
));
}
if matches!(self.shutdown_timeout_ms, Some(0)) {
return Err(OrchestratorError::ConfigLoad(
"Configuration error: `lifecycle_integration.shutdown_timeout_ms` must be at least 1"
.to_string(),
));
}
Ok(())
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
pub struct JudgeCommandsConfig {
#[serde(default)]
pub parallel_dependency: Option<ParallelDependencyJudgeConfig>,
}
impl JudgeCommandsConfig {
fn merge(&mut self, other: Self) {
let Self {
parallel_dependency,
} = other;
overwrite_if_some(&mut self.parallel_dependency, parallel_dependency);
}
pub fn validate(&self) -> Result<()> {
match self.parallel_dependency.as_ref() {
Some(judge) => judge.validate(),
None => Ok(()),
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct ParallelDependencyJudgeConfig {
pub command: Vec<String>,
pub model: String,
#[serde(default)]
pub timeout_ms: Option<u64>,
#[serde(default)]
pub max_input_bytes: Option<usize>,
#[serde(default)]
pub max_output_bytes: Option<usize>,
#[serde(default)]
pub yes_threshold: Option<f64>,
#[serde(default)]
pub mode: Option<String>,
#[serde(default)]
pub evaluation: Option<JudgeEvaluationConfig>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
pub struct JudgeEvaluationConfig {
#[serde(default)]
pub enabled: bool,
#[serde(default)]
pub retention_days: Option<u32>,
#[serde(default)]
pub max_total_bytes: Option<u64>,
}
impl JudgeEvaluationConfig {
fn validate(&self, parent_field: &str) -> Result<()> {
let field = format!("{parent_field}.evaluation");
if let Some(retention_days) = self.retention_days {
if retention_days == 0 || retention_days > MAX_JUDGE_EVALUATION_RETENTION_DAYS {
return Err(OrchestratorError::ConfigLoad(format!(
"Configuration error: `{field}.retention_days` must be between 1 and \
{MAX_JUDGE_EVALUATION_RETENTION_DAYS} (got {retention_days})"
)));
}
}
if let Some(max_total_bytes) = self.max_total_bytes {
if max_total_bytes == 0 || max_total_bytes > MAX_JUDGE_EVALUATION_TOTAL_BYTES {
return Err(OrchestratorError::ConfigLoad(format!(
"Configuration error: `{field}.max_total_bytes` must be between 1 and \
{MAX_JUDGE_EVALUATION_TOTAL_BYTES} (got {max_total_bytes})"
)));
}
}
if !self.enabled {
return Ok(());
}
for (name, present) in [
("retention_days", self.retention_days.is_some()),
("max_total_bytes", self.max_total_bytes.is_some()),
] {
if !present {
return Err(OrchestratorError::ConfigLoad(format!(
"Configuration error: `{field}.{name}` is required when \
`{field}.enabled` is true; evaluation storage is never unbounded"
)));
}
}
Ok(())
}
pub fn bounds(&self) -> Option<(u32, u64)> {
if !self.enabled {
return None;
}
Some((self.retention_days?, self.max_total_bytes?))
}
}
impl ParallelDependencyJudgeConfig {
pub fn timeout_ms(&self) -> u64 {
self.timeout_ms.unwrap_or(DEFAULT_JUDGE_TIMEOUT_MS)
}
pub fn max_input_bytes(&self) -> usize {
self.max_input_bytes
.unwrap_or(DEFAULT_JUDGE_MAX_INPUT_BYTES)
}
pub fn max_output_bytes(&self) -> usize {
self.max_output_bytes
.unwrap_or(DEFAULT_JUDGE_MAX_OUTPUT_BYTES)
}
pub fn yes_threshold(&self) -> f64 {
self.yes_threshold.unwrap_or(DEFAULT_JUDGE_YES_THRESHOLD)
}
pub fn mode(&self) -> &str {
self.mode.as_deref().unwrap_or(JUDGE_MODE_SHADOW)
}
pub fn evaluation_bounds(&self) -> Option<(u32, u64)> {
self.evaluation.as_ref().and_then(|e| e.bounds())
}
pub fn validate(&self) -> Result<()> {
const FIELD: &str = "judge_commands.parallel_dependency";
if self.command.is_empty() {
return Err(OrchestratorError::ConfigLoad(format!(
"Configuration error: `{FIELD}.command` must be a non-empty argv array, for \
example [\"jev\", \"run\", \"-\"]"
)));
}
for (index, element) in self.command.iter().enumerate() {
if element.trim().is_empty() {
return Err(OrchestratorError::ConfigLoad(format!(
"Configuration error: `{FIELD}.command[{index}]` must not be empty"
)));
}
if element.contains('\0') {
return Err(OrchestratorError::ConfigLoad(format!(
"Configuration error: `{FIELD}.command[{index}]` must not contain a NUL byte"
)));
}
}
let model = self.model.trim();
if model.is_empty() {
return Err(OrchestratorError::ConfigLoad(format!(
"Configuration error: `{FIELD}.model` must be a non-empty pinned model identity"
)));
}
if model.ends_with("-latest") {
return Err(OrchestratorError::ConfigLoad(format!(
"Configuration error: `{FIELD}.model` must be a pinned concrete model, not the \
alias `{model}`; the response parser requires the returned model to equal the \
requested one"
)));
}
if let Some(timeout_ms) = self.timeout_ms {
if !(MIN_JUDGE_TIMEOUT_MS..=MAX_JUDGE_TIMEOUT_MS).contains(&timeout_ms) {
return Err(OrchestratorError::ConfigLoad(format!(
"Configuration error: `{FIELD}.timeout_ms` must be between \
{MIN_JUDGE_TIMEOUT_MS} and {MAX_JUDGE_TIMEOUT_MS} (got {timeout_ms})"
)));
}
}
for (field, value) in [
("max_input_bytes", self.max_input_bytes),
("max_output_bytes", self.max_output_bytes),
] {
if let Some(value) = value {
if !(MIN_JUDGE_BYTE_LIMIT..=MAX_JUDGE_BYTE_LIMIT).contains(&value) {
return Err(OrchestratorError::ConfigLoad(format!(
"Configuration error: `{FIELD}.{field}` must be between \
{MIN_JUDGE_BYTE_LIMIT} and {MAX_JUDGE_BYTE_LIMIT} (got {value})"
)));
}
}
}
if let Some(threshold) = self.yes_threshold {
if !threshold.is_finite()
|| !(MIN_JUDGE_YES_THRESHOLD..=MAX_JUDGE_YES_THRESHOLD).contains(&threshold)
{
return Err(OrchestratorError::ConfigLoad(format!(
"Configuration error: `{FIELD}.yes_threshold` must be a finite probability \
between {MIN_JUDGE_YES_THRESHOLD} and {MAX_JUDGE_YES_THRESHOLD} \
(got {threshold})"
)));
}
}
if let Some(mode) = self.mode.as_deref() {
if mode != JUDGE_MODE_SHADOW {
return Err(OrchestratorError::ConfigLoad(format!(
"Configuration error: `{FIELD}.mode` accepts only \
`{JUDGE_MODE_SHADOW}` (got `{mode}`); a judge-authoritative or \
fail-closed mode is not specified yet"
)));
}
}
if let Some(evaluation) = self.evaluation.as_ref() {
evaluation.validate(FIELD)?;
}
Ok(())
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)]
#[serde(rename_all = "snake_case")]
pub enum AcceptancePromptMode {
#[default]
Full,
ContextOnly,
}
fn overwrite_if_some<T>(target: &mut Option<T>, source: Option<T>) {
if source.is_some() {
*target = source;
}
}
fn merge_hooks_config(target: &mut Option<HooksConfig>, source: Option<HooksConfig>) {
match (target.as_mut(), source) {
(Some(target_hooks), Some(source_hooks)) => target_hooks.merge(source_hooks),
(None, Some(source_hooks)) => *target = Some(source_hooks),
(_, None) => {}
}
}
impl OrchestratorConfig {
#[allow(dead_code)]
pub fn new() -> Self {
Self::default()
}
pub fn merge(&mut self, other: Self) {
let Self {
apply_command,
apply_escalation_command,
apply_stall_diagnose_command,
archive_command,
apply_skill,
archive_skill,
analyze_skill,
accept_skill,
rejecting_skill,
cleanup_review_skill,
resolve_skill,
analyze_command,
acceptance_command,
acceptance_escalation_command,
acceptance_escalation,
apply_prompt,
apply_append_prompt,
acceptance_prompt,
acceptance_append_prompt,
acceptance_prompt_mode,
archive_prompt,
archive_append_prompt,
analyze_append_prompt,
resolve_append_prompt,
envs,
hooks,
logging,
stall_detection,
error_circuit_breaker,
completion_check_delay_ms,
completion_check_max_retries,
max_iterations,
max_concurrent_workspaces,
workspace_base_dir,
state_base_dir,
resolve_command,
use_llm_analysis,
vcs_backend,
propose_command,
worktree_command,
command_queue_stagger_delay_ms,
command_queue_max_retries,
command_queue_retry_delay_ms,
command_queue_retry_patterns,
command_queue_retry_if_duration_under_secs,
acceptance_max_continues,
command_inactivity_timeout_secs,
command_inactivity_kill_grace_secs,
command_inactivity_timeout_max_retries,
command_max_runtime_secs,
acceptance_max_runtime_secs,
stream_json_textify,
command_strict_process_cleanup,
lifecycle_integration,
judge_commands,
} = other;
overwrite_if_some(&mut self.apply_command, apply_command);
overwrite_if_some(&mut self.apply_escalation_command, apply_escalation_command);
overwrite_if_some(
&mut self.apply_stall_diagnose_command,
apply_stall_diagnose_command,
);
overwrite_if_some(&mut self.archive_command, archive_command);
overwrite_if_some(&mut self.apply_skill, apply_skill);
overwrite_if_some(&mut self.archive_skill, archive_skill);
overwrite_if_some(&mut self.analyze_skill, analyze_skill);
overwrite_if_some(&mut self.accept_skill, accept_skill);
overwrite_if_some(&mut self.rejecting_skill, rejecting_skill);
overwrite_if_some(&mut self.cleanup_review_skill, cleanup_review_skill);
overwrite_if_some(&mut self.resolve_skill, resolve_skill);
overwrite_if_some(&mut self.analyze_command, analyze_command);
overwrite_if_some(&mut self.acceptance_command, acceptance_command);
overwrite_if_some(
&mut self.acceptance_escalation_command,
acceptance_escalation_command,
);
match (self.acceptance_escalation.as_mut(), acceptance_escalation) {
(Some(target), Some(source)) => target.merge(source),
(None, Some(source)) => self.acceptance_escalation = Some(source),
(_, None) => {}
}
overwrite_if_some(&mut self.resolve_command, resolve_command);
overwrite_if_some(&mut self.apply_prompt, apply_prompt);
overwrite_if_some(&mut self.apply_append_prompt, apply_append_prompt);
overwrite_if_some(&mut self.acceptance_prompt, acceptance_prompt);
overwrite_if_some(&mut self.acceptance_append_prompt, acceptance_append_prompt);
overwrite_if_some(&mut self.archive_prompt, archive_prompt);
overwrite_if_some(&mut self.archive_append_prompt, archive_append_prompt);
overwrite_if_some(&mut self.analyze_append_prompt, analyze_append_prompt);
overwrite_if_some(&mut self.resolve_append_prompt, resolve_append_prompt);
overwrite_if_some(&mut self.acceptance_prompt_mode, acceptance_prompt_mode);
if let Some(envs) = envs {
self.envs.get_or_insert_with(HashMap::new).extend(envs);
}
merge_hooks_config(&mut self.hooks, hooks);
overwrite_if_some(&mut self.logging, logging);
overwrite_if_some(&mut self.stall_detection, stall_detection);
overwrite_if_some(&mut self.error_circuit_breaker, error_circuit_breaker);
overwrite_if_some(
&mut self.completion_check_delay_ms,
completion_check_delay_ms,
);
overwrite_if_some(
&mut self.completion_check_max_retries,
completion_check_max_retries,
);
overwrite_if_some(&mut self.max_iterations, max_iterations);
overwrite_if_some(
&mut self.max_concurrent_workspaces,
max_concurrent_workspaces,
);
overwrite_if_some(&mut self.workspace_base_dir, workspace_base_dir);
overwrite_if_some(&mut self.state_base_dir, state_base_dir);
overwrite_if_some(&mut self.use_llm_analysis, use_llm_analysis);
overwrite_if_some(&mut self.vcs_backend, vcs_backend);
overwrite_if_some(&mut self.propose_command, propose_command);
overwrite_if_some(&mut self.worktree_command, worktree_command);
overwrite_if_some(
&mut self.command_queue_stagger_delay_ms,
command_queue_stagger_delay_ms,
);
overwrite_if_some(
&mut self.command_queue_max_retries,
command_queue_max_retries,
);
overwrite_if_some(
&mut self.command_queue_retry_delay_ms,
command_queue_retry_delay_ms,
);
overwrite_if_some(
&mut self.command_queue_retry_patterns,
command_queue_retry_patterns,
);
overwrite_if_some(
&mut self.command_queue_retry_if_duration_under_secs,
command_queue_retry_if_duration_under_secs,
);
overwrite_if_some(&mut self.acceptance_max_continues, acceptance_max_continues);
overwrite_if_some(
&mut self.command_inactivity_timeout_secs,
command_inactivity_timeout_secs,
);
overwrite_if_some(
&mut self.command_inactivity_kill_grace_secs,
command_inactivity_kill_grace_secs,
);
overwrite_if_some(
&mut self.command_inactivity_timeout_max_retries,
command_inactivity_timeout_max_retries,
);
overwrite_if_some(&mut self.command_max_runtime_secs, command_max_runtime_secs);
overwrite_if_some(
&mut self.acceptance_max_runtime_secs,
acceptance_max_runtime_secs,
);
overwrite_if_some(&mut self.stream_json_textify, stream_json_textify);
overwrite_if_some(
&mut self.command_strict_process_cleanup,
command_strict_process_cleanup,
);
overwrite_if_some(&mut self.lifecycle_integration, lifecycle_integration);
match (self.judge_commands.as_mut(), judge_commands) {
(Some(target), Some(source)) => target.merge(source),
(None, Some(source)) => self.judge_commands = Some(source),
(_, None) => {}
}
}
pub fn get_parallel_dependency_judge(&self) -> Option<&ParallelDependencyJudgeConfig> {
self.judge_commands
.as_ref()
.and_then(|judges| judges.parallel_dependency.as_ref())
}
pub fn validate_judge_commands(&self) -> Result<()> {
match self.judge_commands.as_ref() {
Some(judges) => judges.validate(),
None => Ok(()),
}
}
pub fn get_lifecycle_integration(&self) -> Option<&LifecycleIntegrationConfig> {
self.lifecycle_integration.as_ref()
}
pub fn validate_lifecycle_integration(&self) -> Result<()> {
match self.lifecycle_integration.as_ref() {
Some(integration) => integration.validate(),
None => Ok(()),
}
}
pub fn get_apply_command(&self) -> Result<&str> {
self.apply_command
.as_deref()
.ok_or_else(|| OrchestratorError::ConfigLoad("Missing required config: apply_command. Please set it in .cflx.jsonc or global config.".to_string()))
}
pub fn get_apply_escalation_command(&self) -> Option<&str> {
self.apply_escalation_command.as_deref()
}
pub fn get_apply_stall_diagnose_command(&self) -> Option<&str> {
self.apply_stall_diagnose_command.as_deref()
}
pub fn get_archive_command(&self) -> Result<&str> {
self.archive_command
.as_deref()
.ok_or_else(|| OrchestratorError::ConfigLoad("Missing required config: archive_command. Please set it in .cflx.jsonc or global config.".to_string()))
}
pub fn get_analyze_command(&self) -> Result<&str> {
self.analyze_command
.as_deref()
.ok_or_else(|| OrchestratorError::ConfigLoad("Missing required config: analyze_command. Please set it in .cflx.jsonc or global config.".to_string()))
}
fn configured_operation_skill<'a>(
configured: Option<&'a str>,
default: &'static str,
) -> &'a str {
configured.unwrap_or(default)
}
pub fn get_analyze_skill(&self) -> &str {
Self::configured_operation_skill(self.analyze_skill.as_deref(), DEFAULT_ANALYZE_SKILL)
}
pub fn get_apply_skill(&self) -> &str {
Self::configured_operation_skill(self.apply_skill.as_deref(), DEFAULT_APPLY_SKILL)
}
pub fn get_rejecting_skill(&self) -> &str {
Self::configured_operation_skill(self.rejecting_skill.as_deref(), DEFAULT_REJECTING_SKILL)
}
pub fn get_cleanup_review_skill(&self) -> &str {
Self::configured_operation_skill(
self.cleanup_review_skill.as_deref(),
DEFAULT_CLEANUP_REVIEW_SKILL,
)
}
pub fn get_accept_skill(&self) -> &str {
Self::configured_operation_skill(self.accept_skill.as_deref(), DEFAULT_ACCEPT_SKILL)
}
pub fn get_archive_skill(&self) -> &str {
Self::configured_operation_skill(self.archive_skill.as_deref(), DEFAULT_ARCHIVE_SKILL)
}
pub fn get_resolve_skill(&self) -> &str {
Self::configured_operation_skill(self.resolve_skill.as_deref(), DEFAULT_RESOLVE_SKILL)
}
pub fn get_apply_prompt(&self) -> &str {
self.apply_prompt.as_deref().unwrap_or(DEFAULT_APPLY_PROMPT)
}
pub fn get_archive_prompt(&self) -> &str {
self.archive_prompt
.as_deref()
.unwrap_or(DEFAULT_ARCHIVE_PROMPT)
}
pub fn get_apply_append_prompt(&self) -> Option<&str> {
self.apply_append_prompt.as_deref()
}
pub fn get_acceptance_append_prompt(&self) -> Option<&str> {
self.acceptance_append_prompt.as_deref()
}
pub fn get_archive_append_prompt(&self) -> Option<&str> {
self.archive_append_prompt.as_deref()
}
pub fn get_analyze_append_prompt(&self) -> Option<&str> {
self.analyze_append_prompt.as_deref()
}
pub fn get_resolve_append_prompt(&self) -> Option<&str> {
self.resolve_append_prompt.as_deref()
}
pub fn get_acceptance_command(&self) -> Result<&str> {
self.acceptance_command
.as_deref()
.ok_or_else(|| OrchestratorError::ConfigLoad("Missing required config: acceptance_command. Please set it in .cflx.jsonc or global config.".to_string()))
}
pub fn get_acceptance_escalation_command(&self) -> Option<&str> {
self.acceptance_escalation_command.as_deref()
}
pub fn get_acceptance_escalation(&self) -> AcceptanceEscalationConfig {
self.acceptance_escalation.clone().unwrap_or_default()
}
pub fn validate_acceptance_escalation(&self) -> Result<()> {
match self.acceptance_escalation.as_ref() {
Some(policy) => policy.validate(),
None => Ok(()),
}
}
pub fn get_acceptance_prompt(&self) -> &str {
self.acceptance_prompt
.as_deref()
.unwrap_or(DEFAULT_ACCEPTANCE_PROMPT)
}
pub fn get_acceptance_prompt_mode(&self) -> AcceptancePromptMode {
self.acceptance_prompt_mode.clone().unwrap_or_default()
}
pub fn get_command_envs(&self) -> HashMap<String, String> {
self.envs
.as_ref()
.map(|envs| {
envs.iter()
.map(|(key, value)| (key.clone(), expand::expand_env_value(value)))
.collect()
})
.unwrap_or_default()
}
pub fn get_hooks(&self) -> HooksConfig {
self.hooks.clone().unwrap_or_default()
}
pub fn get_logging(&self) -> LoggingConfig {
self.logging.clone().unwrap_or_default()
}
pub fn get_stall_detection(&self) -> StallDetectionConfig {
self.stall_detection.clone().unwrap_or_default()
}
pub fn get_max_iterations(&self) -> u32 {
self.max_iterations.unwrap_or(DEFAULT_MAX_ITERATIONS)
}
pub fn get_max_concurrent_workspaces(&self) -> usize {
self.max_concurrent_workspaces
.unwrap_or(DEFAULT_MAX_CONCURRENT_WORKSPACES)
}
pub fn get_workspace_base_dir(&self) -> Option<&str> {
self.workspace_base_dir.as_deref().filter(|s| !s.is_empty())
}
pub fn get_state_base_dir(&self) -> Option<&str> {
self.state_base_dir.as_deref().filter(|s| !s.is_empty())
}
pub fn get_resolve_command(&self) -> Result<&str> {
self.resolve_command
.as_deref()
.ok_or_else(|| OrchestratorError::ConfigLoad("Missing required config: resolve_command. Please set it in .cflx.jsonc or global config.".to_string()))
}
pub fn use_llm_analysis(&self) -> bool {
self.use_llm_analysis.unwrap_or(true)
}
pub fn get_vcs_backend(&self) -> VcsBackend {
self.vcs_backend.unwrap_or(VcsBackend::Auto)
}
#[allow(dead_code)]
pub fn get_propose_command(&self) -> Option<&str> {
self.propose_command.as_deref()
}
pub fn get_worktree_command(&self) -> Option<&str> {
self.worktree_command.as_deref()
}
#[allow(dead_code)]
pub fn expand_proposal(template: &str, proposal: &str) -> String {
expand::expand_proposal(template, proposal)
}
pub fn expand_worktree_command(template: &str, workspace_dir: &str, repo_root: &str) -> String {
expand::expand_worktree_command(template, workspace_dir, repo_root)
}
#[allow(dead_code)]
pub fn expand_conflict_files(template: &str, conflict_files: &str) -> String {
expand::expand_conflict_files(template, conflict_files)
}
pub fn get_acceptance_max_continues(&self) -> u32 {
self.acceptance_max_continues
.unwrap_or(defaults::DEFAULT_ACCEPTANCE_MAX_CONTINUES)
}
pub fn get_command_inactivity_timeout_secs(&self) -> u64 {
self.command_inactivity_timeout_secs
.unwrap_or(defaults::DEFAULT_COMMAND_INACTIVITY_TIMEOUT_SECS)
}
pub fn get_command_inactivity_kill_grace_secs(&self) -> u64 {
self.command_inactivity_kill_grace_secs
.unwrap_or(defaults::DEFAULT_COMMAND_INACTIVITY_KILL_GRACE_SECS)
}
pub fn get_command_inactivity_timeout_max_retries(&self) -> u32 {
self.command_inactivity_timeout_max_retries
.unwrap_or(defaults::DEFAULT_COMMAND_INACTIVITY_TIMEOUT_MAX_RETRIES)
}
pub fn get_command_max_runtime_secs(&self) -> u64 {
self.command_max_runtime_secs
.unwrap_or(defaults::DEFAULT_COMMAND_MAX_RUNTIME_SECS)
}
pub fn get_acceptance_max_runtime_secs(&self) -> u64 {
self.acceptance_max_runtime_secs
.unwrap_or(defaults::DEFAULT_ACCEPTANCE_MAX_RUNTIME_SECS)
}
pub fn validate_acceptance_max_runtime_secs(&self) -> Result<()> {
let Some(configured) = self.acceptance_max_runtime_secs else {
return Ok(());
};
let min = defaults::MIN_ACCEPTANCE_MAX_RUNTIME_SECS;
let max = defaults::MAX_ACCEPTANCE_MAX_RUNTIME_SECS;
if (min..=max).contains(&configured) {
return Ok(());
}
Err(OrchestratorError::ConfigLoad(format!(
"Configuration error: `acceptance_max_runtime_secs` must be between {min} and {max} \
seconds (got {configured}). Acceptance cannot disable its absolute runtime limit, so \
`0` is not accepted; use `command_max_runtime_secs` for the common command budget."
)))
}
pub fn get_stream_json_textify(&self) -> bool {
self.stream_json_textify
.unwrap_or(defaults::DEFAULT_STREAM_JSON_TEXTIFY)
}
pub fn get_command_strict_process_cleanup(&self) -> bool {
self.command_strict_process_cleanup
.unwrap_or(defaults::DEFAULT_COMMAND_STRICT_PROCESS_CLEANUP)
}
pub fn expand_change_id(template: &str, change_id: &str) -> String {
expand::expand_change_id(template, change_id)
}
pub fn expand_prompt(template: &str, prompt: &str) -> String {
expand::expand_prompt(template, prompt)
}
fn validate_operation_skill_value(field: &str, value: Option<&str>) -> Result<()> {
let Some(value) = value else {
return Ok(());
};
if value.trim().is_empty() {
return Err(OrchestratorError::ConfigLoad(format!(
"Configuration error: `{field}` must not be empty"
)));
}
if value.contains('\n') || value.contains('\r') {
return Err(OrchestratorError::ConfigLoad(format!(
"Configuration error: `{field}` must not contain newline characters"
)));
}
Ok(())
}
pub fn validate_operation_skills(&self) -> Result<()> {
Self::validate_operation_skill_value("analyze_skill", self.analyze_skill.as_deref())?;
Self::validate_operation_skill_value("apply_skill", self.apply_skill.as_deref())?;
Self::validate_operation_skill_value("rejecting_skill", self.rejecting_skill.as_deref())?;
Self::validate_operation_skill_value(
"cleanup_review_skill",
self.cleanup_review_skill.as_deref(),
)?;
Self::validate_operation_skill_value("accept_skill", self.accept_skill.as_deref())?;
Self::validate_operation_skill_value("archive_skill", self.archive_skill.as_deref())?;
Self::validate_operation_skill_value("resolve_skill", self.resolve_skill.as_deref())?;
Ok(())
}
pub fn validate_required_commands(&self) -> Result<()> {
self.get_stall_detection().validate()?;
self.validate_operation_skills()?;
self.validate_lifecycle_integration()?;
self.validate_acceptance_max_runtime_secs()?;
self.validate_acceptance_escalation()?;
self.validate_judge_commands()?;
let mut missing = Vec::new();
if self.apply_command.is_none() {
missing.push("apply_command");
}
if self.archive_command.is_none() {
missing.push("archive_command");
}
if self.analyze_command.is_none() {
missing.push("analyze_command");
}
if self.acceptance_command.is_none() {
missing.push("acceptance_command");
}
if self.resolve_command.is_none() {
missing.push("resolve_command");
}
if !missing.is_empty() {
return Err(OrchestratorError::ConfigLoad(format!(
"Missing required config: {}. Please set them in .cflx.jsonc or global config.",
missing.join(", ")
)));
}
Ok(())
}
}
#[cfg(test)]
mod lifecycle_integration_config_tests {
use super::*;
fn enabled_config(command: Vec<&str>) -> LifecycleIntegrationConfig {
LifecycleIntegrationConfig {
command: command.into_iter().map(|s| s.to_string()).collect(),
..Default::default()
}
}
#[test]
fn lifecycle_integration_is_absent_by_default() {
let config = OrchestratorConfig::default();
assert!(config.get_lifecycle_integration().is_none());
assert!(config.validate_lifecycle_integration().is_ok());
}
#[test]
fn configured_argv_enables_adapter() {
let integration = enabled_config(vec!["my-adapter", "--flag"]);
assert!(integration.is_enabled());
assert!(integration.validate().is_ok());
}
#[test]
fn explicit_disable_wins_over_configured_argv() {
let integration = LifecycleIntegrationConfig {
enabled: Some(false),
..enabled_config(vec!["my-adapter"])
};
assert!(!integration.is_enabled());
assert!(integration.validate().is_ok());
}
#[test]
fn empty_argv_is_rejected_with_actionable_diagnostic() {
let integration = LifecycleIntegrationConfig {
enabled: Some(true),
..Default::default()
};
let err = integration
.validate()
.expect_err("empty argv must be rejected");
let message = err.to_string();
assert!(
message.contains("lifecycle_integration.command"),
"diagnostic must name the offending key: {message}"
);
assert!(
message.contains("non-empty argv array"),
"diagnostic must explain the expected shape: {message}"
);
}
#[test]
fn blank_executable_is_rejected() {
let integration = enabled_config(vec![" "]);
let err = integration
.validate()
.expect_err("blank executable must be rejected");
assert!(err.to_string().contains("lifecycle_integration.command[0]"));
}
#[test]
fn zero_bounds_are_rejected() {
for integration in [
LifecycleIntegrationConfig {
queue_capacity: Some(0),
..enabled_config(vec!["my-adapter"])
},
LifecycleIntegrationConfig {
write_timeout_ms: Some(0),
..enabled_config(vec!["my-adapter"])
},
LifecycleIntegrationConfig {
shutdown_timeout_ms: Some(0),
..enabled_config(vec!["my-adapter"])
},
] {
assert!(
integration.validate().is_err(),
"zero bound must be rejected: {integration:?}"
);
}
}
#[test]
fn bounded_defaults_are_applied_when_unset() {
let integration = enabled_config(vec!["my-adapter"]);
assert_eq!(
integration.queue_capacity(),
DEFAULT_LIFECYCLE_QUEUE_CAPACITY
);
assert_eq!(
integration.write_timeout_ms(),
DEFAULT_LIFECYCLE_WRITE_TIMEOUT_MS
);
assert_eq!(
integration.shutdown_timeout_ms(),
DEFAULT_LIFECYCLE_SHUTDOWN_TIMEOUT_MS
);
}
#[test]
fn required_command_validation_surfaces_lifecycle_errors() {
let config = OrchestratorConfig {
apply_command: Some("apply".to_string()),
archive_command: Some("archive".to_string()),
analyze_command: Some("analyze".to_string()),
acceptance_command: Some("accept".to_string()),
resolve_command: Some("resolve".to_string()),
lifecycle_integration: Some(LifecycleIntegrationConfig {
enabled: Some(true),
..Default::default()
}),
..Default::default()
};
let err = config
.validate_required_commands()
.expect_err("invalid lifecycle integration must fail required-command validation");
assert!(err.to_string().contains("lifecycle_integration.command"));
}
#[test]
fn merge_overwrites_lifecycle_integration_and_preserves_unset() {
let mut base = OrchestratorConfig {
lifecycle_integration: Some(enabled_config(vec!["base-adapter"])),
..Default::default()
};
base.merge(OrchestratorConfig::default());
assert_eq!(
base.get_lifecycle_integration().map(|i| i.command.clone()),
Some(vec!["base-adapter".to_string()])
);
base.merge(OrchestratorConfig {
lifecycle_integration: Some(enabled_config(vec!["project-adapter"])),
..Default::default()
});
assert_eq!(
base.get_lifecycle_integration().map(|i| i.command.clone()),
Some(vec!["project-adapter".to_string()])
);
}
#[test]
fn existing_configuration_without_lifecycle_key_still_parses() {
let parsed: OrchestratorConfig =
serde_json::from_str(r#"{"apply_command": "apply", "archive_command": "archive"}"#)
.expect("legacy configuration must remain parseable");
assert!(parsed.get_lifecycle_integration().is_none());
}
#[test]
fn lifecycle_integration_parses_from_configuration_document() {
let parsed: OrchestratorConfig = serde_json::from_str(
r#"{"lifecycle_integration": {"command": ["adapter", "--pane"], "shutdown_timeout_ms": 500}}"#,
)
.expect("lifecycle integration must parse");
let integration = parsed
.get_lifecycle_integration()
.expect("integration should be present");
assert!(integration.is_enabled());
assert_eq!(integration.command, vec!["adapter", "--pane"]);
assert_eq!(integration.shutdown_timeout_ms(), 500);
}
}
#[cfg(test)]
mod judge_commands_config_tests {
use super::*;
fn valid_judge() -> ParallelDependencyJudgeConfig {
ParallelDependencyJudgeConfig {
command: vec!["jev".to_string(), "run".to_string(), "-".to_string()],
model: "jev-1.13.0".to_string(),
timeout_ms: None,
max_input_bytes: None,
max_output_bytes: None,
yes_threshold: None,
mode: None,
evaluation: None,
}
}
fn judge_config(judge: ParallelDependencyJudgeConfig) -> OrchestratorConfig {
OrchestratorConfig {
judge_commands: Some(JudgeCommandsConfig {
parallel_dependency: Some(judge),
}),
..Default::default()
}
}
#[test]
fn judge_commands_are_absent_by_default() {
let config = OrchestratorConfig::default();
assert!(config.judge_commands.is_none());
assert!(config.get_parallel_dependency_judge().is_none());
assert!(config.validate_judge_commands().is_ok());
}
#[test]
fn existing_configuration_without_judge_commands_still_parses() {
let parsed: OrchestratorConfig = serde_json::from_str(
r#"{"apply_command": "apply", "archive_command": "archive", "analyze_command": "analyze", "acceptance_command": "accept", "resolve_command": "resolve"}"#,
)
.expect("configuration without judge_commands must remain parseable");
assert!(parsed.get_parallel_dependency_judge().is_none());
assert!(parsed.validate_required_commands().is_ok());
}
#[test]
fn valid_argv_model_and_bounds_load_with_defaults() {
let parsed: OrchestratorConfig = serde_json::from_str(
r#"{"judge_commands": {"parallel_dependency": {"command": ["jev", "run", "-"], "model": "jev-1.13.0"}}}"#,
)
.expect("judge command must parse");
let judge = parsed
.get_parallel_dependency_judge()
.expect("entry should be present");
assert_eq!(judge.command, vec!["jev", "run", "-"]);
assert_eq!(judge.model, "jev-1.13.0");
assert_eq!(judge.timeout_ms(), DEFAULT_JUDGE_TIMEOUT_MS);
assert_eq!(judge.max_input_bytes(), DEFAULT_JUDGE_MAX_INPUT_BYTES);
assert_eq!(judge.max_output_bytes(), DEFAULT_JUDGE_MAX_OUTPUT_BYTES);
assert_eq!(judge.yes_threshold(), DEFAULT_JUDGE_YES_THRESHOLD);
assert!(judge.validate().is_ok());
}
#[test]
fn explicit_bounds_override_defaults() {
let judge = ParallelDependencyJudgeConfig {
timeout_ms: Some(5_000),
max_input_bytes: Some(2_048),
max_output_bytes: Some(4_096),
yes_threshold: Some(0.6),
..valid_judge()
};
assert!(judge.validate().is_ok());
assert_eq!(judge.timeout_ms(), 5_000);
assert_eq!(judge.max_input_bytes(), 2_048);
assert_eq!(judge.max_output_bytes(), 4_096);
assert_eq!(judge.yes_threshold(), 0.6);
}
#[test]
fn invalid_entries_are_rejected_naming_the_field() {
let cases: Vec<(ParallelDependencyJudgeConfig, &str)> = vec![
(
ParallelDependencyJudgeConfig {
command: Vec::new(),
..valid_judge()
},
"parallel_dependency.command",
),
(
ParallelDependencyJudgeConfig {
command: vec!["jev".to_string(), " ".to_string()],
..valid_judge()
},
"parallel_dependency.command[1]",
),
(
ParallelDependencyJudgeConfig {
command: vec!["je\0v".to_string()],
..valid_judge()
},
"parallel_dependency.command[0]",
),
(
ParallelDependencyJudgeConfig {
model: " ".to_string(),
..valid_judge()
},
"parallel_dependency.model",
),
(
ParallelDependencyJudgeConfig {
model: "jev-latest".to_string(),
..valid_judge()
},
"parallel_dependency.model",
),
(
ParallelDependencyJudgeConfig {
timeout_ms: Some(0),
..valid_judge()
},
"parallel_dependency.timeout_ms",
),
(
ParallelDependencyJudgeConfig {
timeout_ms: Some(MAX_JUDGE_TIMEOUT_MS + 1),
..valid_judge()
},
"parallel_dependency.timeout_ms",
),
(
ParallelDependencyJudgeConfig {
max_input_bytes: Some(0),
..valid_judge()
},
"parallel_dependency.max_input_bytes",
),
(
ParallelDependencyJudgeConfig {
max_output_bytes: Some(MAX_JUDGE_BYTE_LIMIT + 1),
..valid_judge()
},
"parallel_dependency.max_output_bytes",
),
(
ParallelDependencyJudgeConfig {
yes_threshold: Some(f64::NAN),
..valid_judge()
},
"parallel_dependency.yes_threshold",
),
(
ParallelDependencyJudgeConfig {
yes_threshold: Some(0.4),
..valid_judge()
},
"parallel_dependency.yes_threshold",
),
(
ParallelDependencyJudgeConfig {
yes_threshold: Some(1.5),
..valid_judge()
},
"parallel_dependency.yes_threshold",
),
];
for (judge, expected_field) in cases {
let message = judge
.validate()
.expect_err("invalid judge entry must be rejected")
.to_string();
assert!(
message.contains(expected_field),
"diagnostic must name `{expected_field}`: {message}"
);
}
}
#[test]
fn required_command_validation_surfaces_judge_errors() {
let config = OrchestratorConfig {
apply_command: Some("apply".to_string()),
archive_command: Some("archive".to_string()),
analyze_command: Some("analyze".to_string()),
acceptance_command: Some("accept".to_string()),
resolve_command: Some("resolve".to_string()),
..judge_config(ParallelDependencyJudgeConfig {
command: Vec::new(),
..valid_judge()
})
};
let message = config
.validate_required_commands()
.expect_err("an invalid judge entry must fail startup validation")
.to_string();
assert!(message.contains("judge_commands.parallel_dependency.command"));
}
#[test]
fn higher_priority_entry_replaces_the_complete_trust_boundary() {
let mut lower = judge_config(ParallelDependencyJudgeConfig {
command: vec!["global-judge".to_string()],
model: "global-1.0.0".to_string(),
timeout_ms: Some(1_000),
max_input_bytes: Some(2_048),
max_output_bytes: Some(4_096),
yes_threshold: Some(0.95),
mode: Some(JUDGE_MODE_SHADOW.to_string()),
evaluation: Some(JudgeEvaluationConfig {
enabled: true,
retention_days: Some(90),
max_total_bytes: Some(1_048_576),
}),
});
lower.merge(judge_config(ParallelDependencyJudgeConfig {
command: vec!["project-judge".to_string()],
model: "project-2.0.0".to_string(),
timeout_ms: None,
max_input_bytes: None,
max_output_bytes: None,
yes_threshold: None,
mode: None,
evaluation: None,
}));
let merged = lower
.get_parallel_dependency_judge()
.expect("entry should be present");
assert_eq!(merged.command, vec!["project-judge".to_string()]);
assert_eq!(merged.model, "project-2.0.0");
assert_eq!(merged.timeout_ms, None);
assert_eq!(merged.timeout_ms(), DEFAULT_JUDGE_TIMEOUT_MS);
assert_eq!(merged.max_input_bytes(), DEFAULT_JUDGE_MAX_INPUT_BYTES);
assert_eq!(merged.max_output_bytes(), DEFAULT_JUDGE_MAX_OUTPUT_BYTES);
assert_eq!(merged.yes_threshold(), DEFAULT_JUDGE_YES_THRESHOLD);
assert_eq!(merged.mode, None);
assert_eq!(merged.mode(), JUDGE_MODE_SHADOW);
assert_eq!(merged.evaluation, None);
assert_eq!(merged.evaluation_bounds(), None);
}
#[test]
fn omitted_entry_inherits_the_lower_priority_entry() {
let mut lower = judge_config(valid_judge());
lower.merge(OrchestratorConfig::default());
assert_eq!(
lower
.get_parallel_dependency_judge()
.map(|j| j.model.clone()),
Some("jev-1.13.0".to_string())
);
lower.merge(OrchestratorConfig {
judge_commands: Some(JudgeCommandsConfig::default()),
..Default::default()
});
assert_eq!(
lower
.get_parallel_dependency_judge()
.map(|j| j.model.clone()),
Some("jev-1.13.0".to_string())
);
}
#[test]
fn judge_entry_is_inherited_when_no_lower_priority_entry_exists() {
let mut lower = OrchestratorConfig::default();
lower.merge(judge_config(valid_judge()));
assert_eq!(
lower
.get_parallel_dependency_judge()
.map(|j| j.command.clone()),
Some(vec!["jev".to_string(), "run".to_string(), "-".to_string()])
);
}
mod judge_evaluation_config {
use super::*;
fn parse(json: &str) -> OrchestratorConfig {
serde_json::from_str(json).expect("configuration must parse")
}
fn judge_from(json: &str) -> ParallelDependencyJudgeConfig {
parse(json)
.get_parallel_dependency_judge()
.expect("judge entry must be present")
.clone()
}
fn with_evaluation(evaluation: &str) -> String {
format!(
r#"{{"judge_commands": {{"parallel_dependency": {{
"command": ["jev", "run", "-"], "model": "jev-1.13.0",
"evaluation": {evaluation}}}}}}}"#
)
}
#[test]
fn an_omitted_mode_normalizes_to_shadow() {
let judge = judge_from(
r#"{"judge_commands": {"parallel_dependency": {"command": ["jev"], "model": "m-1"}}}"#,
);
assert_eq!(judge.mode, None);
assert_eq!(judge.mode(), JUDGE_MODE_SHADOW);
assert!(judge.validate().is_ok());
}
#[test]
fn shadow_is_the_only_accepted_mode() {
let judge = judge_from(
r#"{"judge_commands": {"parallel_dependency": {"command": ["jev"], "model": "m-1", "mode": "shadow"}}}"#,
);
assert_eq!(judge.mode(), JUDGE_MODE_SHADOW);
assert!(judge.validate().is_ok());
for rejected in ["prefer", "only", "authoritative", "Shadow", ""] {
let judge = ParallelDependencyJudgeConfig {
mode: Some(rejected.to_string()),
..valid_judge()
};
let message = judge
.validate()
.expect_err(&format!("`{rejected}` must not be an accepted mode"))
.to_string();
assert!(
message.contains("judge_commands.parallel_dependency.mode"),
"the diagnostic must name the field: {message}"
);
}
}
#[test]
fn evaluation_is_absent_by_default_and_records_nothing() {
let judge = judge_from(
r#"{"judge_commands": {"parallel_dependency": {"command": ["jev"], "model": "m-1"}}}"#,
);
assert_eq!(judge.evaluation, None);
assert_eq!(judge.evaluation_bounds(), None);
assert!(judge.validate().is_ok());
}
#[test]
fn an_enabled_policy_exposes_its_bounds() {
let judge = judge_from(&with_evaluation(
r#"{"enabled": true, "retention_days": 30, "max_total_bytes": 104857600}"#,
));
assert!(judge.validate().is_ok());
assert_eq!(judge.evaluation_bounds(), Some((30, 104_857_600)));
}
#[test]
fn a_disabled_policy_exposes_no_bounds_even_when_it_carries_them() {
let judge = judge_from(&with_evaluation(
r#"{"enabled": false, "retention_days": 30, "max_total_bytes": 1024}"#,
));
assert!(judge.validate().is_ok(), "the bounds are still validated");
assert_eq!(judge.evaluation_bounds(), None);
let omitted_enabled = judge_from(&with_evaluation(
r#"{"retention_days": 30, "max_total_bytes": 1024}"#,
));
assert_eq!(
omitted_enabled.evaluation_bounds(),
None,
"an omitted `enabled` is disabled"
);
}
#[test]
fn an_enabled_policy_requires_both_bounds() {
for (label, evaluation) in [
("neither bound", r#"{"enabled": true}"#),
(
"only retention",
r#"{"enabled": true, "retention_days": 30}"#,
),
(
"only bytes",
r#"{"enabled": true, "max_total_bytes": 1024}"#,
),
] {
let message = parse(&with_evaluation(evaluation))
.validate_judge_commands()
.expect_err(&format!("{label} must fail configuration loading"))
.to_string();
assert!(
message.contains("judge_commands.parallel_dependency.evaluation"),
"the diagnostic must name the field for {label}: {message}"
);
assert!(
message.contains("required"),
"and say what is missing for {label}: {message}"
);
}
}
#[test]
fn bounds_are_rejected_at_zero_and_above_their_maximums() {
let cases = [
("retention_days", "0", "retention_days"),
("retention_days", "3651", "retention_days"),
("max_total_bytes", "0", "max_total_bytes"),
("max_total_bytes", "10737418241", "max_total_bytes"),
];
for (field, value, expected_field) in cases {
let other = if field == "retention_days" {
r#""max_total_bytes": 1024"#
} else {
r#""retention_days": 30"#
};
let message = parse(&with_evaluation(&format!(
r#"{{"enabled": true, "{field}": {value}, {other}}}"#
)))
.validate_judge_commands()
.expect_err(&format!("{field} = {value} must be rejected"))
.to_string();
assert!(
message.contains(expected_field),
"the diagnostic must name `{expected_field}`: {message}"
);
}
let judge = judge_from(&with_evaluation(
r#"{"enabled": true, "retention_days": 3650, "max_total_bytes": 10737418240}"#,
));
assert!(judge.validate().is_ok());
assert_eq!(
judge.evaluation_bounds(),
Some((
MAX_JUDGE_EVALUATION_RETENTION_DAYS,
MAX_JUDGE_EVALUATION_TOTAL_BYTES
))
);
}
#[test]
fn out_of_range_bounds_are_rejected_even_while_disabled() {
let message = parse(&with_evaluation(
r#"{"enabled": false, "retention_days": 99999, "max_total_bytes": 1024}"#,
))
.validate_judge_commands()
.expect_err("an out-of-range bound must fail even while disabled")
.to_string();
assert!(message.contains("retention_days"), "{message}");
}
#[test]
fn a_project_layer_cannot_form_a_hybrid_evaluation_policy() {
let mut lower = judge_config(ParallelDependencyJudgeConfig {
mode: Some(JUDGE_MODE_SHADOW.to_string()),
evaluation: Some(JudgeEvaluationConfig {
enabled: true,
retention_days: Some(365),
max_total_bytes: Some(1_048_576),
}),
..valid_judge()
});
lower.merge(judge_config(ParallelDependencyJudgeConfig {
command: vec!["project-judge".to_string()],
model: "project-2.0.0".to_string(),
..valid_judge()
}));
let merged = lower
.get_parallel_dependency_judge()
.expect("entry must be present");
assert_eq!(merged.model, "project-2.0.0");
assert_eq!(merged.evaluation, None);
assert_eq!(
merged.evaluation_bounds(),
None,
"recording must not be inherited across a replaced entry"
);
let mut lower = judge_config(valid_judge());
lower.merge(judge_config(ParallelDependencyJudgeConfig {
evaluation: Some(JudgeEvaluationConfig {
enabled: true,
retention_days: Some(7),
max_total_bytes: Some(2048),
}),
..valid_judge()
}));
assert_eq!(
lower
.get_parallel_dependency_judge()
.and_then(|j| j.evaluation_bounds()),
Some((7, 2048))
);
}
#[test]
fn a_configuration_without_judge_commands_still_has_no_evaluation_policy() {
let config = parse(r#"{"apply_command": "echo apply"}"#);
assert!(config.get_parallel_dependency_judge().is_none());
assert!(config.validate_judge_commands().is_ok());
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn operation_skill_accessors_return_defaults_when_unset() {
let config = OrchestratorConfig::default();
assert_eq!(config.get_analyze_skill(), DEFAULT_ANALYZE_SKILL);
assert_eq!(config.get_apply_skill(), DEFAULT_APPLY_SKILL);
assert_eq!(config.get_rejecting_skill(), DEFAULT_REJECTING_SKILL);
assert_eq!(
config.get_cleanup_review_skill(),
DEFAULT_CLEANUP_REVIEW_SKILL
);
assert_eq!(config.get_accept_skill(), DEFAULT_ACCEPT_SKILL);
assert_eq!(config.get_archive_skill(), DEFAULT_ARCHIVE_SKILL);
assert_eq!(config.get_resolve_skill(), DEFAULT_RESOLVE_SKILL);
}
#[test]
fn operation_skill_accessors_return_configured_values() {
let config = OrchestratorConfig {
analyze_skill: Some("team-analyze".to_string()),
apply_skill: Some("team-apply".to_string()),
rejecting_skill: Some("team-rejecting".to_string()),
cleanup_review_skill: Some("team-cleanup-review".to_string()),
accept_skill: Some("cflx-accept-with-speca".to_string()),
archive_skill: Some("team-archive".to_string()),
resolve_skill: Some("team-resolve".to_string()),
..Default::default()
};
assert_eq!(config.get_analyze_skill(), "team-analyze");
assert_eq!(config.get_apply_skill(), "team-apply");
assert_eq!(config.get_rejecting_skill(), "team-rejecting");
assert_eq!(config.get_cleanup_review_skill(), "team-cleanup-review");
assert_eq!(config.get_accept_skill(), "cflx-accept-with-speca");
assert_eq!(config.get_archive_skill(), "team-archive");
assert_eq!(config.get_resolve_skill(), "team-resolve");
}
#[test]
fn operation_skill_merge_preserves_lower_precedence_when_omitted() {
let mut lower = OrchestratorConfig {
accept_skill: Some("cflx-accept-with-speca".to_string()),
resolve_skill: Some("team-resolve".to_string()),
..Default::default()
};
lower.merge(OrchestratorConfig::default());
assert_eq!(lower.get_accept_skill(), "cflx-accept-with-speca");
assert_eq!(lower.get_resolve_skill(), "team-resolve");
}
#[test]
fn operation_skill_merge_overrides_accept_and_resolve() {
let mut config = OrchestratorConfig {
accept_skill: Some("lower-accept".to_string()),
resolve_skill: Some("lower-resolve".to_string()),
..Default::default()
};
let higher = OrchestratorConfig {
accept_skill: Some("cflx-accept-with-speca".to_string()),
resolve_skill: Some("team-resolve".to_string()),
..Default::default()
};
config.merge(higher);
assert_eq!(config.get_accept_skill(), "cflx-accept-with-speca");
assert_eq!(config.get_resolve_skill(), "team-resolve");
}
#[test]
fn operation_skill_validation_rejects_empty_or_newline_values() {
let empty = OrchestratorConfig {
accept_skill: Some(" ".to_string()),
..Default::default()
};
assert!(empty.validate_operation_skills().is_err());
let newline = OrchestratorConfig {
resolve_skill: Some("team-resolve\nload skills: other".to_string()),
..Default::default()
};
assert!(newline.validate_operation_skills().is_err());
}
}