use std::cell::RefCell;
use std::collections::{BTreeMap, VecDeque};
use serde::{Deserialize, Serialize};
use super::action::{Action, DecisionClass, Finding};
use super::driver::ChunkState;
use super::envelope::{DecisionEnvelope, DecisionTier};
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "trigger", rename_all = "snake_case")]
pub enum DecisionTrigger {
SpecReady,
ChunkCommitted {
chunk_id: String,
},
VerifyReport {
report_id: String,
findings: Vec<Finding>,
},
CircuitBreakerTripped {
reason: String,
},
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct DecisionContext {
pub run_id: String,
pub plan_rev: u32,
pub intent_rev: u32,
pub chunks: BTreeMap<String, ChunkState>,
pub trigger: DecisionTrigger,
}
pub trait Orchestrator {
fn decide(&self, ctx: &DecisionContext) -> Vec<(Action, DecisionEnvelope)>;
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CoordinatorProposal {
pub action: Action,
pub reason: String,
pub input_artifacts: Vec<String>,
}
pub trait Coordinator {
fn coordinate(&self, ctx: &DecisionContext) -> Vec<CoordinatorProposal>;
fn model(&self) -> String;
fn prompt_version(&self) -> String;
fn actor(&self) -> String {
"coordinator".to_string()
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DeciderVerdict {
pub action: Action,
pub reason: String,
pub input_artifacts: Vec<String>,
}
pub trait Decider {
fn decide_consequential(
&self,
ctx: &DecisionContext,
proposed: &CoordinatorProposal,
) -> DeciderVerdict;
fn model(&self) -> String;
fn prompt_version(&self) -> String;
fn actor(&self) -> String {
"decider".to_string()
}
}
pub struct TieredOrchestrator<C, D> {
coordinator: C,
decider: D,
}
impl<C: Coordinator, D: Decider> TieredOrchestrator<C, D> {
pub fn new(coordinator: C, decider: D) -> Self {
Self {
coordinator,
decider,
}
}
}
impl<C: Coordinator, D: Decider> Orchestrator for TieredOrchestrator<C, D> {
fn decide(&self, ctx: &DecisionContext) -> Vec<(Action, DecisionEnvelope)> {
self.coordinator
.coordinate(ctx)
.into_iter()
.map(|proposal| route_proposal(&self.coordinator, &self.decider, ctx, proposal))
.collect()
}
}
pub fn route_proposal(
coordinator: &dyn Coordinator,
decider: &dyn Decider,
ctx: &DecisionContext,
proposal: CoordinatorProposal,
) -> (Action, DecisionEnvelope) {
match proposal.action.decision_class() {
DecisionClass::Routine => {
let envelope = DecisionEnvelope {
actor: coordinator.actor(),
input_artifacts: proposal.input_artifacts,
reason: proposal.reason,
decision_tier: DecisionTier::Coordinator,
model: coordinator.model(),
prompt_version: coordinator.prompt_version(),
};
(proposal.action, envelope)
}
DecisionClass::Consequential => {
let verdict = decider.decide_consequential(ctx, &proposal);
let envelope = DecisionEnvelope {
actor: decider.actor(),
input_artifacts: verdict.input_artifacts,
reason: verdict.reason,
decision_tier: DecisionTier::Decider,
model: decider.model(),
prompt_version: decider.prompt_version(),
};
debug_assert!(
envelope.validate_for(&verdict.action).is_ok(),
"route_proposal produced a tier-invariant violation for {}",
verdict.action.name()
);
(verdict.action, envelope)
}
}
}
pub struct ScriptedCoordinator {
script: RefCell<VecDeque<Vec<CoordinatorProposal>>>,
model: String,
prompt_version: String,
}
impl ScriptedCoordinator {
pub fn new(batches: Vec<Vec<CoordinatorProposal>>) -> Self {
Self {
script: RefCell::new(batches.into()),
model: "stub-coordinator".to_string(),
prompt_version: "stub-v1".to_string(),
}
}
}
impl Coordinator for ScriptedCoordinator {
fn coordinate(&self, _ctx: &DecisionContext) -> Vec<CoordinatorProposal> {
self.script.borrow_mut().pop_front().unwrap_or_default()
}
fn model(&self) -> String {
self.model.clone()
}
fn prompt_version(&self) -> String {
self.prompt_version.clone()
}
}
pub struct ScriptedDecider {
script: RefCell<VecDeque<DeciderVerdict>>,
model: String,
prompt_version: String,
}
impl ScriptedDecider {
pub fn new(verdicts: Vec<DeciderVerdict>) -> Self {
Self {
script: RefCell::new(verdicts.into()),
model: "stub-decider".to_string(),
prompt_version: "stub-v1".to_string(),
}
}
pub fn confirming() -> Self {
Self::new(Vec::new())
}
}
impl Decider for ScriptedDecider {
fn decide_consequential(
&self,
_ctx: &DecisionContext,
proposed: &CoordinatorProposal,
) -> DeciderVerdict {
self.script
.borrow_mut()
.pop_front()
.unwrap_or_else(|| DeciderVerdict {
action: proposed.action.clone(),
reason: format!("decider confirmed: {}", proposed.reason),
input_artifacts: proposed.input_artifacts.clone(),
})
}
fn model(&self) -> String {
self.model.clone()
}
fn prompt_version(&self) -> String {
self.prompt_version.clone()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::pipeline::action::SpinoffScope;
fn ctx() -> DecisionContext {
DecisionContext {
run_id: "run1".into(),
plan_rev: 1,
intent_rev: 1,
chunks: BTreeMap::new(),
trigger: DecisionTrigger::SpecReady,
}
}
fn proposal(action: Action) -> CoordinatorProposal {
CoordinatorProposal {
action,
reason: "because".into(),
input_artifacts: vec!["plan:1".into()],
}
}
#[test]
fn routine_proposal_is_stamped_coordinator() {
let coord = ScriptedCoordinator::new(vec![vec![proposal(Action::AcceptChunk {
chunk_id: "c1".into(),
})]]);
let orch = TieredOrchestrator::new(coord, ScriptedDecider::confirming());
let decisions = orch.decide(&ctx());
assert_eq!(decisions.len(), 1);
let (action, env) = &decisions[0];
assert_eq!(action.name(), "accept_chunk");
assert_eq!(env.decision_tier, DecisionTier::Coordinator);
assert!(env.validate_for(action).is_ok());
}
#[test]
fn consequential_proposal_is_deferred_and_stamped_decider() {
let coord = ScriptedCoordinator::new(vec![vec![proposal(Action::DeclareConverged)]]);
let orch = TieredOrchestrator::new(coord, ScriptedDecider::confirming());
let decisions = orch.decide(&ctx());
assert_eq!(decisions.len(), 1);
let (action, env) = &decisions[0];
assert_eq!(action.name(), "declare_converged");
assert_eq!(env.decision_tier, DecisionTier::Decider);
assert_eq!(env.model, "stub-decider");
assert!(env.validate_for(action).is_ok());
}
#[test]
fn decider_may_override_the_proposed_action() {
let coord = ScriptedCoordinator::new(vec![vec![proposal(Action::DeclareConverged)]]);
let decider = ScriptedDecider::new(vec![DeciderVerdict {
action: Action::Escalate {
reason: "not actually done".into(),
},
reason: "intent not met".into(),
input_artifacts: vec!["intent:1".into()],
}]);
let orch = TieredOrchestrator::new(coord, decider);
let decisions = orch.decide(&ctx());
let (action, env) = &decisions[0];
assert_eq!(action.name(), "escalate");
assert_eq!(env.decision_tier, DecisionTier::Decider);
assert_eq!(env.reason, "intent not met");
}
#[test]
fn nontrivial_spinoff_is_deferred_but_trivial_stays_coordinator() {
let coord = ScriptedCoordinator::new(vec![vec![
proposal(Action::ProposeSpinoff {
title: "trivial".into(),
kind: "improvement".into(),
rationale: "r".into(),
scope: SpinoffScope::Trivial,
}),
proposal(Action::ProposeSpinoff {
title: "big".into(),
kind: "refactor".into(),
rationale: "r".into(),
scope: SpinoffScope::Substantial,
}),
]]);
let orch = TieredOrchestrator::new(coord, ScriptedDecider::confirming());
let decisions = orch.decide(&ctx());
assert_eq!(decisions.len(), 2);
assert_eq!(decisions[0].1.decision_tier, DecisionTier::Coordinator);
assert_eq!(decisions[1].1.decision_tier, DecisionTier::Decider);
}
#[test]
fn exhausted_coordinator_script_yields_no_decisions() {
let coord = ScriptedCoordinator::new(vec![]);
let orch = TieredOrchestrator::new(coord, ScriptedDecider::confirming());
assert!(orch.decide(&ctx()).is_empty());
}
}