use std::time::{Duration, Instant};
use serde::{Deserialize, Serialize};
use crate::state::NodeId;
#[derive(Debug, Clone)]
pub struct CognitivePipeline {
pub(crate) operators: Vec<CognitiveOperator>,
pub(crate) budget: Option<Duration>,
pub(crate) explain: bool,
pub(crate) mode: ExecutionMode,
pub(crate) context: Option<PipelineContext>,
}
impl CognitivePipeline {
pub fn new() -> Self {
Self {
operators: Vec::new(),
budget: None,
explain: false,
mode: ExecutionMode::Lazy,
context: None,
}
}
pub fn with_context(mut self, ctx: PipelineContext) -> Self {
self.context = Some(ctx);
self
}
pub fn with_budget(mut self, budget: Duration) -> Self {
self.budget = Some(budget);
self.mode = ExecutionMode::Budgeted;
self
}
pub fn with_mode(mut self, mode: ExecutionMode) -> Self {
self.mode = mode;
self
}
pub fn attend(mut self, seeds: Vec<NodeId>) -> Self {
self.operators.push(CognitiveOperator::Attend(AttendOp {
seeds,
max_hops: 2,
decay: 0.5,
}));
self
}
pub fn attend_with(mut self, seeds: Vec<NodeId>, max_hops: u32, decay: f64) -> Self {
self.operators.push(CognitiveOperator::Attend(AttendOp {
seeds,
max_hops,
decay,
}));
self
}
pub fn recall(mut self, top_k: usize) -> Self {
self.operators.push(CognitiveOperator::Recall(RecallOp {
top_k,
query: None,
domain: None,
}));
self
}
pub fn recall_query(mut self, query: String, top_k: usize, domain: Option<String>) -> Self {
self.operators.push(CognitiveOperator::Recall(RecallOp {
top_k,
query: Some(query),
domain,
}));
self
}
pub fn believe(mut self, evidence: EvidenceInput) -> Self {
self.operators.push(CognitiveOperator::Believe(BelieveOp {
evidence,
}));
self
}
pub fn project(mut self, horizon: ProjectionHorizon) -> Self {
self.operators.push(CognitiveOperator::Project(ProjectOp {
horizon,
include_causal: true,
}));
self
}
pub fn compare(mut self, candidates: Vec<CandidateAction>) -> Self {
self.operators.push(CognitiveOperator::Compare(CompareOp {
candidates,
apply_personality: true,
}));
self
}
pub fn constrain(mut self, policies: Vec<PolicyConstraint>) -> Self {
self.operators.push(CognitiveOperator::Constrain(ConstrainOp {
policies,
}));
self
}
pub fn anticipate(mut self, horizon_secs: f64) -> Self {
self.operators.push(CognitiveOperator::Anticipate(AnticipateOp {
horizon_secs,
}));
self
}
pub fn plan(mut self, goal_id: NodeId, max_depth: u32) -> Self {
self.operators.push(CognitiveOperator::Plan(PlanOp {
goal_id,
max_depth,
}));
self
}
pub fn assess(mut self) -> Self {
self.operators.push(CognitiveOperator::Assess);
self
}
pub fn coherence_check(mut self) -> Self {
self.operators.push(CognitiveOperator::CoherenceCheck);
self
}
pub fn explain(mut self) -> Self {
self.explain = true;
self
}
pub fn len(&self) -> usize {
self.operators.len()
}
pub fn is_empty(&self) -> bool {
self.operators.is_empty()
}
}
impl Default for CognitivePipeline {
fn default() -> Self {
Self::new()
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum CognitiveOperator {
Attend(AttendOp),
Recall(RecallOp),
Believe(BelieveOp),
Project(ProjectOp),
Compare(CompareOp),
Constrain(ConstrainOp),
Anticipate(AnticipateOp),
Plan(PlanOp),
Assess,
CoherenceCheck,
}
impl CognitiveOperator {
pub fn name(&self) -> &'static str {
match self {
Self::Attend(_) => "attend",
Self::Recall(_) => "recall",
Self::Believe(_) => "believe",
Self::Project(_) => "project",
Self::Compare(_) => "compare",
Self::Constrain(_) => "constrain",
Self::Anticipate(_) => "anticipate",
Self::Plan(_) => "plan",
Self::Assess => "assess",
Self::CoherenceCheck => "coherence_check",
}
}
pub fn priority(&self) -> u8 {
match self {
Self::Attend(_) => 10, Self::Recall(_) => 9, Self::Believe(_) => 8, Self::Compare(_) => 7, Self::Constrain(_) => 7, Self::Plan(_) => 6, Self::Project(_) => 5, Self::Anticipate(_) => 4, Self::Assess => 3, Self::CoherenceCheck => 2, }
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AttendOp {
pub seeds: Vec<NodeId>,
pub max_hops: u32,
pub decay: f64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RecallOp {
pub top_k: usize,
pub query: Option<String>,
pub domain: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct BelieveOp {
pub evidence: EvidenceInput,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct EvidenceInput {
pub target: Option<NodeId>,
pub observation: String,
pub direction: f64,
pub strength: f64,
pub source: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ProjectOp {
pub horizon: ProjectionHorizon,
pub include_causal: bool,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
pub enum ProjectionHorizon {
OneStep,
ShortTerm,
MediumTerm,
LongTerm,
}
impl ProjectionHorizon {
pub fn seconds(&self) -> f64 {
match self {
Self::OneStep => 60.0,
Self::ShortTerm => 600.0, Self::MediumTerm => 7200.0, Self::LongTerm => 86400.0, }
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CompareOp {
pub candidates: Vec<CandidateAction>,
pub apply_personality: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CandidateAction {
pub description: String,
pub action_kind: String,
pub confidence: f64,
pub properties: crate::personality_bias::ActionProperties,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ConstrainOp {
pub policies: Vec<PolicyConstraint>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PolicyConstraint {
pub name: String,
pub kind: ConstraintKind,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum ConstraintKind {
MinConfidence(f64),
MaxRisk(f64),
MetaCogThreshold(f64),
RequireConsent(String),
BlockActionKind(String),
QuietHours { start_hour: u8, end_hour: u8 },
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AnticipateOp {
pub horizon_secs: f64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PlanOp {
pub goal_id: NodeId,
pub max_depth: u32,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum ExecutionMode {
Lazy,
Budgeted,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PipelineContext {
pub trigger: String,
pub user_input: Option<String>,
pub mentioned_entities: Vec<NodeId>,
pub timestamp: f64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PipelineResult {
pub steps: Vec<StepResult>,
pub operators_executed: usize,
pub operators_skipped: usize,
pub elapsed_ms: u64,
pub budget_exhausted: bool,
pub explanation: Option<ExplanationTrace>,
pub status: PipelineStatus,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum PipelineStatus {
Complete,
Partial,
Failed,
Empty,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct StepResult {
pub operator: String,
pub success: bool,
pub elapsed_ms: u64,
pub output: StepOutput,
pub trace: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum StepOutput {
Attend {
nodes_activated: usize,
top_activated: Vec<(NodeId, f64)>,
},
Recall {
memories_retrieved: usize,
top_matches: Vec<RecallMatch>,
},
Believe {
beliefs_updated: usize,
confidence_delta: f64,
},
Project {
predictions: Vec<Prediction>,
},
Compare {
ranked: Vec<RankedCandidate>,
},
Constrain {
passed: usize,
filtered_out: usize,
violations: Vec<ConstraintViolation>,
},
Anticipate {
items: Vec<AnticipatedItem>,
},
Plan {
plan_found: bool,
steps: usize,
score: f64,
},
Assess {
overall_confidence: f64,
coverage_gaps: usize,
},
Coherence {
score: f64,
conflicts: usize,
stale_nodes: usize,
},
Skipped {
reason: String,
},
Error {
message: String,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RecallMatch {
pub text: String,
pub similarity: f64,
pub memory_type: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Prediction {
pub description: String,
pub probability: f64,
pub valence: f64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RankedCandidate {
pub description: String,
pub score: f64,
pub personality_bias: f64,
pub rank: usize,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ConstraintViolation {
pub constraint_name: String,
pub candidate: String,
pub reason: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AnticipatedItem {
pub description: String,
pub expected_at: f64,
pub confidence: f64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ExplanationTrace {
pub steps: Vec<ExplanationStep>,
pub summary: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ExplanationStep {
pub step_number: usize,
pub operator: String,
pub description: String,
pub key_insight: Option<String>,
}
pub fn execute_pipeline(
pipeline: &CognitivePipeline,
executor: &dyn PipelineExecutor,
) -> PipelineResult {
if pipeline.is_empty() {
return PipelineResult {
steps: Vec::new(),
operators_executed: 0,
operators_skipped: 0,
elapsed_ms: 0,
budget_exhausted: false,
explanation: None,
status: PipelineStatus::Empty,
};
}
let start = Instant::now();
let mut steps = Vec::with_capacity(pipeline.operators.len());
let mut executed = 0usize;
let mut skipped = 0usize;
let mut budget_exhausted = false;
let mut any_failed = false;
let execution_order: Vec<usize> = if pipeline.mode == ExecutionMode::Budgeted {
let mut indices: Vec<usize> = (0..pipeline.operators.len()).collect();
indices.sort_by(|a, b| {
pipeline.operators[*b].priority().cmp(&pipeline.operators[*a].priority())
});
indices
} else {
(0..pipeline.operators.len()).collect()
};
for &idx in &execution_order {
if let Some(budget) = pipeline.budget {
if start.elapsed() >= budget {
budget_exhausted = true;
skipped += execution_order.len() - executed;
for &remaining_idx in &execution_order[executed..] {
let _ = remaining_idx; }
break;
}
}
let op = &pipeline.operators[idx];
let step_start = Instant::now();
let output = executor.execute_operator(op, &pipeline.context);
let step_elapsed = step_start.elapsed().as_millis() as u64;
let success = !matches!(output, StepOutput::Error { .. });
if !success {
any_failed = true;
}
let trace = if pipeline.explain {
Some(format_step_trace(op, &output))
} else {
None
};
steps.push(StepResult {
operator: op.name().to_string(),
success,
elapsed_ms: step_elapsed,
output,
trace,
});
executed += 1;
}
let total_elapsed = start.elapsed().as_millis() as u64;
let status = if any_failed {
PipelineStatus::Failed
} else if budget_exhausted {
PipelineStatus::Partial
} else {
PipelineStatus::Complete
};
let explanation = if pipeline.explain {
Some(build_explanation(&steps))
} else {
None
};
PipelineResult {
steps,
operators_executed: executed,
operators_skipped: skipped,
elapsed_ms: total_elapsed,
budget_exhausted,
explanation,
status,
}
}
pub trait PipelineExecutor {
fn execute_operator(
&self,
operator: &CognitiveOperator,
context: &Option<PipelineContext>,
) -> StepOutput;
}
fn format_step_trace(op: &CognitiveOperator, output: &StepOutput) -> String {
match (op, output) {
(CognitiveOperator::Attend(a), StepOutput::Attend { nodes_activated, .. }) => {
format!("Spread activation from {} seeds → {} nodes activated",
a.seeds.len(), nodes_activated)
}
(CognitiveOperator::Recall(r), StepOutput::Recall { memories_retrieved, .. }) => {
format!("Retrieved {} memories (top_k={})",
memories_retrieved, r.top_k)
}
(CognitiveOperator::Believe(_), StepOutput::Believe { beliefs_updated, confidence_delta, .. }) => {
format!("Updated {} beliefs (Δconf={:+.3})", beliefs_updated, confidence_delta)
}
(CognitiveOperator::Project(p), StepOutput::Project { predictions, .. }) => {
format!("Projected {:?} → {} predictions", p.horizon, predictions.len())
}
(CognitiveOperator::Compare(_), StepOutput::Compare { ranked, .. }) => {
format!("Compared {} candidates → top: {}",
ranked.len(),
ranked.first().map(|r| r.description.as_str()).unwrap_or("none"))
}
(CognitiveOperator::Constrain(_), StepOutput::Constrain { passed, filtered_out, .. }) => {
format!("{} passed, {} filtered by constraints", passed, filtered_out)
}
(CognitiveOperator::Plan(_), StepOutput::Plan { plan_found, steps, score, .. }) => {
if *plan_found {
format!("Plan found: {} steps, score {:.2}", steps, score)
} else {
"No viable plan found".to_string()
}
}
(CognitiveOperator::Assess, StepOutput::Assess { overall_confidence, coverage_gaps, .. }) => {
format!("Meta-cognitive confidence: {:.2}, {} coverage gaps",
overall_confidence, coverage_gaps)
}
(CognitiveOperator::CoherenceCheck, StepOutput::Coherence { score, conflicts, stale_nodes, .. }) => {
format!("Coherence: {:.2}, {} conflicts, {} stale", score, conflicts, stale_nodes)
}
(_, StepOutput::Skipped { reason }) => {
format!("Skipped: {}", reason)
}
(_, StepOutput::Error { message }) => {
format!("Error: {}", message)
}
_ => format!("{}: completed", op.name()),
}
}
fn build_explanation(steps: &[StepResult]) -> ExplanationTrace {
let explanation_steps: Vec<ExplanationStep> = steps.iter().enumerate()
.filter(|(_, s)| s.success)
.map(|(i, s)| {
ExplanationStep {
step_number: i + 1,
operator: s.operator.clone(),
description: s.trace.clone().unwrap_or_else(|| format!("{}: completed", s.operator)),
key_insight: extract_key_insight(&s.output),
}
})
.collect();
let summary = if explanation_steps.is_empty() {
"No reasoning steps executed.".to_string()
} else {
let op_names: Vec<&str> = explanation_steps.iter()
.map(|s| s.operator.as_str())
.collect();
format!("Reasoning pipeline: {} steps ({})", explanation_steps.len(), op_names.join(" → "))
};
ExplanationTrace {
steps: explanation_steps,
summary,
}
}
fn extract_key_insight(output: &StepOutput) -> Option<String> {
match output {
StepOutput::Compare { ranked, .. } if !ranked.is_empty() => {
Some(format!("Best candidate: {} (score {:.2})",
ranked[0].description, ranked[0].score))
}
StepOutput::Assess { overall_confidence, .. } if *overall_confidence < 0.5 => {
Some("Low meta-cognitive confidence — consider escalating".to_string())
}
StepOutput::Coherence { score, .. } if *score < 0.5 => {
Some("Coherence degraded — enforcement may be needed".to_string())
}
StepOutput::Constrain { filtered_out, .. } if *filtered_out > 0 => {
Some(format!("{} candidates filtered by policy constraints", filtered_out))
}
_ => None,
}
}
pub struct PipelinePatterns;
impl PipelinePatterns {
pub fn user_turn(seeds: Vec<NodeId>, candidates: Vec<CandidateAction>) -> CognitivePipeline {
CognitivePipeline::new()
.attend(seeds)
.recall(5)
.compare(candidates)
.constrain(vec![PolicyConstraint {
name: "min_confidence".to_string(),
kind: ConstraintKind::MinConfidence(0.3),
}])
.explain()
}
pub fn proactive(candidates: Vec<CandidateAction>, horizon_secs: f64) -> CognitivePipeline {
CognitivePipeline::new()
.anticipate(horizon_secs)
.assess()
.compare(candidates)
.constrain(vec![
PolicyConstraint {
name: "min_confidence".to_string(),
kind: ConstraintKind::MinConfidence(0.5),
},
PolicyConstraint {
name: "meta_cog".to_string(),
kind: ConstraintKind::MetaCogThreshold(0.4),
},
])
.explain()
}
pub fn deep_reasoning(
seeds: Vec<NodeId>,
goal_id: NodeId,
candidates: Vec<CandidateAction>,
) -> CognitivePipeline {
CognitivePipeline::new()
.attend(seeds)
.recall(10)
.project(ProjectionHorizon::MediumTerm)
.plan(goal_id, 4)
.compare(candidates)
.constrain(vec![PolicyConstraint {
name: "min_confidence".to_string(),
kind: ConstraintKind::MinConfidence(0.4),
}])
.assess()
.explain()
}
pub fn health_check() -> CognitivePipeline {
CognitivePipeline::new()
.assess()
.coherence_check()
.explain()
}
pub fn budgeted(
seeds: Vec<NodeId>,
candidates: Vec<CandidateAction>,
budget_ms: u64,
) -> CognitivePipeline {
CognitivePipeline::new()
.with_budget(Duration::from_millis(budget_ms))
.attend(seeds)
.recall(5)
.compare(candidates)
.constrain(vec![PolicyConstraint {
name: "min_confidence".to_string(),
kind: ConstraintKind::MinConfidence(0.3),
}])
.assess()
.explain()
}
}
pub struct StubExecutor;
impl PipelineExecutor for StubExecutor {
fn execute_operator(
&self,
operator: &CognitiveOperator,
_context: &Option<PipelineContext>,
) -> StepOutput {
match operator {
CognitiveOperator::Attend(a) => StepOutput::Attend {
nodes_activated: a.seeds.len() * 3,
top_activated: a.seeds.iter().map(|id| (*id, 0.8)).collect(),
},
CognitiveOperator::Recall(r) => StepOutput::Recall {
memories_retrieved: r.top_k.min(3),
top_matches: vec![RecallMatch {
text: "Test memory".to_string(),
similarity: 0.85,
memory_type: "episodic".to_string(),
}],
},
CognitiveOperator::Believe(_) => StepOutput::Believe {
beliefs_updated: 1,
confidence_delta: 0.1,
},
CognitiveOperator::Project(_) => StepOutput::Project {
predictions: vec![Prediction {
description: "User will check email".to_string(),
probability: 0.7,
valence: 0.3,
}],
},
CognitiveOperator::Compare(c) => {
let ranked: Vec<RankedCandidate> = c.candidates.iter().enumerate()
.map(|(i, cand)| RankedCandidate {
description: cand.description.clone(),
score: cand.confidence * 0.9,
personality_bias: 0.05,
rank: i + 1,
})
.collect();
StepOutput::Compare { ranked }
}
CognitiveOperator::Constrain(_) => StepOutput::Constrain {
passed: 2,
filtered_out: 1,
violations: vec![ConstraintViolation {
constraint_name: "min_confidence".to_string(),
candidate: "Low confidence action".to_string(),
reason: "Confidence 0.2 < threshold 0.3".to_string(),
}],
},
CognitiveOperator::Anticipate(_) => StepOutput::Anticipate {
items: vec![AnticipatedItem {
description: "Meeting in 30 minutes".to_string(),
expected_at: 1800.0,
confidence: 0.9,
}],
},
CognitiveOperator::Plan(_) => StepOutput::Plan {
plan_found: true,
steps: 3,
score: 0.75,
},
CognitiveOperator::Assess => StepOutput::Assess {
overall_confidence: 0.72,
coverage_gaps: 2,
},
CognitiveOperator::CoherenceCheck => StepOutput::Coherence {
score: 0.85,
conflicts: 0,
stale_nodes: 1,
},
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::state::NodeKind;
fn seed_ids() -> Vec<NodeId> {
vec![
NodeId::new(NodeKind::Entity, 1),
NodeId::new(NodeKind::Entity, 2),
]
}
fn test_candidates() -> Vec<CandidateAction> {
vec![
CandidateAction {
description: "Send notification".to_string(),
action_kind: "notify".to_string(),
confidence: 0.8,
properties: Default::default(),
},
CandidateAction {
description: "Wait and observe".to_string(),
action_kind: "wait".to_string(),
confidence: 0.6,
properties: Default::default(),
},
]
}
#[test]
fn test_empty_pipeline() {
let pipeline = CognitivePipeline::new();
let result = execute_pipeline(&pipeline, &StubExecutor);
assert_eq!(result.status, PipelineStatus::Empty);
assert_eq!(result.operators_executed, 0);
}
#[test]
fn test_simple_pipeline() {
let pipeline = CognitivePipeline::new()
.attend(seed_ids())
.recall(5)
.compare(test_candidates());
assert_eq!(pipeline.len(), 3);
let result = execute_pipeline(&pipeline, &StubExecutor);
assert_eq!(result.status, PipelineStatus::Complete);
assert_eq!(result.operators_executed, 3);
assert_eq!(result.operators_skipped, 0);
}
#[test]
fn test_pipeline_with_explanation() {
let pipeline = CognitivePipeline::new()
.attend(seed_ids())
.recall(5)
.explain();
let result = execute_pipeline(&pipeline, &StubExecutor);
assert!(result.explanation.is_some());
let trace = result.explanation.unwrap();
assert_eq!(trace.steps.len(), 2);
assert!(trace.summary.contains("2 steps"));
}
#[test]
fn test_user_turn_pattern() {
let pipeline = PipelinePatterns::user_turn(seed_ids(), test_candidates());
assert!(pipeline.len() >= 4);
assert!(pipeline.explain);
let result = execute_pipeline(&pipeline, &StubExecutor);
assert_eq!(result.status, PipelineStatus::Complete);
}
#[test]
fn test_health_check_pattern() {
let pipeline = PipelinePatterns::health_check();
assert_eq!(pipeline.len(), 2);
let result = execute_pipeline(&pipeline, &StubExecutor);
assert_eq!(result.status, PipelineStatus::Complete);
let coherence_step = result.steps.iter()
.find(|s| s.operator == "coherence_check")
.unwrap();
assert!(coherence_step.success);
}
#[test]
fn test_deep_reasoning_pattern() {
let goal_id = NodeId::new(NodeKind::Goal, 1);
let pipeline = PipelinePatterns::deep_reasoning(
seed_ids(), goal_id, test_candidates(),
);
assert!(pipeline.len() >= 6);
let result = execute_pipeline(&pipeline, &StubExecutor);
assert_eq!(result.status, PipelineStatus::Complete);
}
#[test]
fn test_operator_priority_ordering() {
assert!(CognitiveOperator::Attend(AttendOp {
seeds: vec![], max_hops: 1, decay: 0.5,
}).priority() > CognitiveOperator::Assess.priority());
assert!(CognitiveOperator::Recall(RecallOp {
top_k: 5, query: None, domain: None,
}).priority() > CognitiveOperator::CoherenceCheck.priority());
}
#[test]
fn test_operator_names() {
assert_eq!(CognitiveOperator::Attend(AttendOp {
seeds: vec![], max_hops: 1, decay: 0.5,
}).name(), "attend");
assert_eq!(CognitiveOperator::Assess.name(), "assess");
assert_eq!(CognitiveOperator::CoherenceCheck.name(), "coherence_check");
}
#[test]
fn test_projection_horizons() {
assert!(ProjectionHorizon::OneStep.seconds() < ProjectionHorizon::ShortTerm.seconds());
assert!(ProjectionHorizon::ShortTerm.seconds() < ProjectionHorizon::MediumTerm.seconds());
assert!(ProjectionHorizon::MediumTerm.seconds() < ProjectionHorizon::LongTerm.seconds());
}
#[test]
fn test_budgeted_pipeline_short_budget() {
let pipeline = CognitivePipeline::new()
.with_budget(Duration::from_millis(500))
.attend(seed_ids())
.recall(5)
.assess()
.coherence_check();
assert_eq!(pipeline.mode, ExecutionMode::Budgeted);
let result = execute_pipeline(&pipeline, &StubExecutor);
assert_eq!(result.status, PipelineStatus::Complete);
assert_eq!(result.operators_executed, 4);
}
#[test]
fn test_pipeline_context() {
let pipeline = CognitivePipeline::new()
.with_context(PipelineContext {
trigger: "user_message".to_string(),
user_input: Some("Hello".to_string()),
mentioned_entities: seed_ids(),
timestamp: 1_000_000.0,
})
.attend(seed_ids());
assert!(pipeline.context.is_some());
assert_eq!(pipeline.context.as_ref().unwrap().trigger, "user_message");
}
#[test]
fn test_evidence_input() {
let evidence = EvidenceInput {
target: Some(NodeId::new(NodeKind::Belief, 1)),
observation: "User confirmed preference".to_string(),
direction: 1.0,
strength: 0.8,
source: "user".to_string(),
};
let pipeline = CognitivePipeline::new()
.believe(evidence);
assert_eq!(pipeline.len(), 1);
let result = execute_pipeline(&pipeline, &StubExecutor);
assert_eq!(result.status, PipelineStatus::Complete);
}
#[test]
fn test_constraint_kinds() {
let constraints = vec![
PolicyConstraint {
name: "confidence".to_string(),
kind: ConstraintKind::MinConfidence(0.5),
},
PolicyConstraint {
name: "risk".to_string(),
kind: ConstraintKind::MaxRisk(0.3),
},
PolicyConstraint {
name: "quiet_hours".to_string(),
kind: ConstraintKind::QuietHours { start_hour: 22, end_hour: 7 },
},
];
let pipeline = CognitivePipeline::new()
.constrain(constraints);
assert_eq!(pipeline.len(), 1);
}
#[test]
fn test_key_insight_extraction() {
let output = StepOutput::Assess {
overall_confidence: 0.3,
coverage_gaps: 5,
};
let insight = extract_key_insight(&output);
assert!(insight.is_some());
assert!(insight.unwrap().contains("escalating"));
}
#[test]
fn test_proactive_pattern() {
let pipeline = PipelinePatterns::proactive(test_candidates(), 3600.0);
assert!(pipeline.len() >= 4);
let result = execute_pipeline(&pipeline, &StubExecutor);
assert_eq!(result.status, PipelineStatus::Complete);
}
}