Skip to main content

ggen_core/ontology/
control_loop.rs

1use serde::{Deserialize, Serialize};
2/// Autonomous Control Loop: Closed-Loop Ontology Evolution
3///
4/// Implements the complete feedback loop:
5/// Observe → Detect → Propose → Validate → Promote → Record → Repeat
6///
7/// Runs autonomously without human intervention in the editing loop.
8use std::sync::Arc;
9use std::time::Duration;
10use tokio::sync::RwLock;
11
12use crate::ontology::delta_proposer::DeltaSigmaProposer;
13#[cfg(test)]
14use crate::ontology::pattern_miner::ObservationSource;
15use crate::ontology::pattern_miner::{MinerConfig, Observation, PatternMiner};
16use crate::ontology::promotion::AtomicSnapshotPromoter;
17use crate::ontology::sigma_runtime::{SigmaReceipt, SigmaRuntime, SigmaSnapshot};
18use crate::ontology::validators::{CompositeValidator, Invariant, ValidationContext};
19
20/// Telemetry for a control loop iteration
21#[derive(Debug, Clone, Serialize, Deserialize)]
22pub struct IterationTelemetry {
23    pub iteration: usize,
24    pub timestamp_ms: u64,
25    pub observation_count: usize,
26    pub patterns_detected: usize,
27    pub proposals_generated: usize,
28    pub proposals_validated: usize,
29    pub proposals_promoted: usize,
30    pub total_duration_ms: u64,
31}
32
33/// Control loop state machine
34#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
35pub enum LoopState {
36    Idle,
37    Observing,
38    Detecting,
39    Proposing,
40    Validating,
41    Promoting,
42    Recording,
43    Error,
44}
45
46/// Autonomous control loop configuration
47#[derive(Debug, Clone)]
48pub struct ControlLoopConfig {
49    /// Interval between iterations (milliseconds)
50    pub iteration_interval_ms: u64,
51
52    /// Max iterations before stopping (None = infinite)
53    pub max_iterations: Option<usize>,
54
55    /// Enable automatic promotion of valid proposals
56    pub auto_promote: bool,
57
58    /// Sector to evolve
59    pub sector: String,
60
61    /// Min confidence to proceed with proposal
62    pub min_proposal_confidence: f64,
63
64    /// Miner configuration
65    pub miner_config: MinerConfig,
66}
67
68impl Default for ControlLoopConfig {
69    fn default() -> Self {
70        Self {
71            iteration_interval_ms: 5000,
72            max_iterations: None,
73            auto_promote: true,
74            sector: "support".to_string(),
75            min_proposal_confidence: 0.75,
76            miner_config: MinerConfig::default(),
77        }
78    }
79}
80
81/// The autonomous control loop
82pub struct AutonomousControlLoop {
83    config: ControlLoopConfig,
84    state: Arc<RwLock<LoopState>>,
85    sigma_runtime: Arc<RwLock<SigmaRuntime>>,
86    promoter: Arc<AtomicSnapshotPromoter>,
87    pattern_miner: Arc<RwLock<PatternMiner>>,
88    proposer: Arc<dyn DeltaSigmaProposer>,
89    validator: Arc<CompositeValidator>,
90    telemetry: Arc<RwLock<Vec<IterationTelemetry>>>,
91}
92
93impl AutonomousControlLoop {
94    pub fn new(
95        config: ControlLoopConfig, initial_snapshot: SigmaSnapshot,
96        proposer: Arc<dyn DeltaSigmaProposer>, validator: Arc<CompositeValidator>,
97    ) -> Self {
98        let sigma_runtime = SigmaRuntime::new(initial_snapshot.clone());
99        let miner_config = config.miner_config.clone();
100
101        Self {
102            config,
103            state: Arc::new(RwLock::new(LoopState::Idle)),
104            sigma_runtime: Arc::new(RwLock::new(sigma_runtime)),
105            promoter: Arc::new(AtomicSnapshotPromoter::new(Arc::new(initial_snapshot))),
106            pattern_miner: Arc::new(RwLock::new(PatternMiner::new(miner_config))),
107            proposer,
108            validator,
109            telemetry: Arc::new(RwLock::new(Vec::new())),
110        }
111    }
112
113    /// Get current loop state
114    pub async fn state(&self) -> LoopState {
115        *self.state.read().await
116    }
117
118    /// Get telemetry
119    pub async fn telemetry(&self) -> Vec<IterationTelemetry> {
120        self.telemetry.read().await.clone()
121    }
122
123    /// Feed observation into the system
124    pub async fn observe(&self, obs: Observation) {
125        let mut miner = self.pattern_miner.write().await;
126        miner.add_observation(obs);
127    }
128
129    /// Run one iteration of the control loop
130    async fn iterate(&self) -> Result<IterationTelemetry, String> {
131        let start_ms = get_time_ms();
132        let mut telemetry = IterationTelemetry {
133            iteration: self.telemetry.read().await.len(),
134            timestamp_ms: start_ms,
135            observation_count: 0,
136            patterns_detected: 0,
137            proposals_generated: 0,
138            proposals_validated: 0,
139            proposals_promoted: 0,
140            total_duration_ms: 0,
141        };
142
143        // 1. OBSERVE: Already done via .observe() calls
144        let miner = self.pattern_miner.read().await;
145        telemetry.observation_count = miner.observation_count();
146        drop(miner);
147
148        // 2. DETECT: Run pattern mining
149        *self.state.write().await = LoopState::Detecting;
150        let mut miner = self.pattern_miner.write().await;
151        let patterns = miner
152            .mine()
153            .map_err(|e| format!("Pattern mining failed: {}", e))?;
154        telemetry.patterns_detected = patterns.len();
155        drop(miner);
156
157        if patterns.is_empty() {
158            *self.state.write().await = LoopState::Recording;
159            return Ok(telemetry);
160        }
161
162        // 3. PROPOSE: Generate ΔΣ² proposals
163        *self.state.write().await = LoopState::Proposing;
164        let current_snapshot = self
165            .promoter
166            .get_current()
167            .map_err(|e| format!("Failed to get current snapshot: {}", e))?;
168        let proposals = self
169            .proposer
170            .propose_deltas(patterns, current_snapshot.snapshot(), &self.config.sector)
171            .await
172            .map_err(|e| format!("Proposal generation failed: {}", e))?;
173
174        telemetry.proposals_generated = proposals.len();
175
176        // Filter by confidence
177        let valid_proposals: Vec<_> = proposals
178            .iter()
179            .filter(|p| p.confidence >= self.config.min_proposal_confidence)
180            .cloned()
181            .collect();
182
183        // 4. VALIDATE: Check invariants (Q)
184        *self.state.write().await = LoopState::Validating;
185        for proposal in &valid_proposals {
186            let current_snap = self
187                .promoter
188                .get_current()
189                .map_err(|e| format!("Failed to get current snapshot: {}", e))?;
190
191            // Apply proposal changes to create new snapshot
192            let mut new_triples = current_snap.snapshot().triples.as_ref().clone();
193
194            // Remove triples (for now, just filter out matching subjects)
195            for triple_pattern in &proposal.triples_to_remove {
196                new_triples.retain(|stmt| !stmt.subject.contains(triple_pattern));
197            }
198
199            // Add new triples
200            for triple_str in &proposal.triples_to_add {
201                // Parse simple triple format: "subject predicate object"
202                let parts: Vec<&str> = triple_str.split_whitespace().collect();
203                if parts.len() >= 3 {
204                    new_triples.push(crate::ontology::sigma_runtime::Statement {
205                        subject: parts[0].to_string(),
206                        predicate: parts[1].to_string(),
207                        object: parts[2].to_string(),
208                        graph: None,
209                    });
210                }
211            }
212
213            let new_snap = SigmaSnapshot::new(
214                Some(current_snap.snapshot().id.clone()),
215                new_triples,
216                format!("{}_updated", current_snap.snapshot().version),
217                "sig_updated".to_string(),
218                current_snap.snapshot().metadata.clone(),
219            );
220
221            let ctx = ValidationContext {
222                proposal: proposal.clone(),
223                current_snapshot: current_snap.snapshot(),
224                expected_new_snapshot: Arc::new(new_snap),
225                sector: self.config.sector.clone(),
226                invariants: vec![
227                    Invariant::NoRetrocausation,
228                    Invariant::TypeSoundness,
229                    Invariant::SLOPreservation,
230                ],
231            };
232
233            let (static_ev, dynamic_ev, perf_ev) = self
234                .validator
235                .validate_all(&ctx)
236                .await
237                .map_err(|e| format!("Validation failed: {}", e))?;
238
239            telemetry.proposals_validated += 1;
240
241            // Check if all validations passed
242            if static_ev.passed && dynamic_ev.passed && perf_ev.passed {
243                // 5. PROMOTE: Move to current
244                *self.state.write().await = LoopState::Promoting;
245
246                if self.config.auto_promote {
247                    // Create the promoted snapshot with applied changes
248                    let current_snap_for_promote = self.promoter.get_current().map_err(|e| {
249                        format!("Failed to get current snapshot for promotion: {}", e)
250                    })?;
251
252                    let mut promoted_triples =
253                        current_snap_for_promote.snapshot().triples.as_ref().clone();
254
255                    // Apply the same changes for promotion
256                    for triple_pattern in &proposal.triples_to_remove {
257                        promoted_triples.retain(|stmt| !stmt.subject.contains(triple_pattern));
258                    }
259
260                    for triple_str in &proposal.triples_to_add {
261                        let parts: Vec<&str> = triple_str.split_whitespace().collect();
262                        if parts.len() >= 3 {
263                            promoted_triples.push(crate::ontology::sigma_runtime::Statement {
264                                subject: parts[0].to_string(),
265                                predicate: parts[1].to_string(),
266                                object: parts[2].to_string(),
267                                graph: None,
268                            });
269                        }
270                    }
271
272                    let promoted_snapshot = SigmaSnapshot::new(
273                        Some(current_snap_for_promote.snapshot().id.clone()),
274                        promoted_triples,
275                        format!(
276                            "{}_v{}",
277                            current_snap_for_promote.snapshot().version,
278                            telemetry.iteration
279                        ),
280                        "promoted_sig".to_string(),
281                        current_snap_for_promote.snapshot().metadata.clone(),
282                    );
283
284                    let _promotion_result = self
285                        .promoter
286                        .promote(Arc::new(promoted_snapshot))
287                        .map_err(|e| format!("Failed to promote snapshot: {}", e))?;
288
289                    telemetry.proposals_promoted += 1;
290
291                    // 6. RECORD: Store receipt
292                    let receipt = SigmaReceipt::new(
293                        Default::default(),
294                        Some(current_snap.snapshot().id.clone()),
295                        format!("Proposal: {}", proposal.id),
296                    );
297
298                    let mut runtime = self.sigma_runtime.write().await;
299                    runtime.record_receipt(receipt);
300                }
301            }
302        }
303
304        // 7. RECORD: Update telemetry
305        *self.state.write().await = LoopState::Recording;
306        telemetry.total_duration_ms = get_time_ms() - start_ms;
307
308        Ok(telemetry)
309    }
310
311    /// Run the control loop
312    pub async fn run(&self) -> Result<(), String> {
313        *self.state.write().await = LoopState::Idle;
314
315        let mut iteration = 0;
316        loop {
317            // Check max iterations
318            if let Some(max) = self.config.max_iterations {
319                if iteration >= max {
320                    break;
321                }
322            }
323
324            // Run iteration
325            match self.iterate().await {
326                Ok(telemetry) => {
327                    self.telemetry.write().await.push(telemetry);
328                }
329                Err(e) => {
330                    *self.state.write().await = LoopState::Error;
331                    return Err(e);
332                }
333            }
334
335            iteration += 1;
336
337            // Wait before next iteration
338            tokio::time::sleep(Duration::from_millis(self.config.iteration_interval_ms)).await;
339        }
340
341        *self.state.write().await = LoopState::Idle;
342        Ok(())
343    }
344
345    /// Run with bounded iterations
346    pub async fn run_bounded(&self, max_iters: usize) -> Result<(), String> {
347        // Note: we can't modify self.config (it's not mutable), so we'll track manually
348
349        for _ in 0..max_iters {
350            match self.iterate().await {
351                Ok(telemetry) => {
352                    self.telemetry.write().await.push(telemetry);
353                }
354                Err(e) => {
355                    *self.state.write().await = LoopState::Error;
356                    return Err(e);
357                }
358            }
359
360            tokio::time::sleep(Duration::from_millis(self.config.iteration_interval_ms)).await;
361        }
362
363        *self.state.write().await = LoopState::Idle;
364        Ok(())
365    }
366
367    /// Get current snapshot
368    pub fn current_snapshot(&self) -> Result<Arc<SigmaSnapshot>, String> {
369        self.promoter
370            .get_current()
371            .map_err(|e| format!("Failed to get current snapshot: {}", e))
372            .map(|guard| guard.snapshot())
373    }
374}
375
376/// Get current time in milliseconds
377fn get_time_ms() -> u64 {
378    use std::time::{SystemTime, UNIX_EPOCH};
379
380    SystemTime::now()
381        .duration_since(UNIX_EPOCH)
382        .unwrap_or_default()
383        .as_millis() as u64
384}
385
386#[cfg(test)]
387mod tests {
388    use super::*;
389    use crate::ontology::delta_proposer::PatternHeuristicProposer;
390    use crate::ontology::validators::{
391        RealDynamicValidator, RealPerformanceValidator, RealStaticValidator,
392    };
393
394    fn create_test_loop() -> AutonomousControlLoop {
395        let snapshot = SigmaSnapshot::new(
396            None,
397            vec![],
398            "1.0.0".to_string(),
399            "sig".to_string(),
400            Default::default(),
401        );
402
403        let proposer: Arc<dyn DeltaSigmaProposer> =
404            Arc::new(PatternHeuristicProposer::new(Default::default()));
405
406        let static_v: Arc<dyn crate::ontology::validators::StaticValidator> =
407            Arc::new(RealStaticValidator::new(vec![]));
408        let dynamic_v: Arc<dyn crate::ontology::validators::DynamicValidator> =
409            Arc::new(RealDynamicValidator::new(3));
410        let perf_v: Arc<dyn crate::ontology::validators::PerformanceValidator> =
411            Arc::new(RealPerformanceValidator::new(10000, 100 * 1024));
412
413        let validator = Arc::new(CompositeValidator::new(static_v, dynamic_v, perf_v));
414
415        let config = ControlLoopConfig {
416            iteration_interval_ms: 100,
417            ..Default::default()
418        };
419
420        AutonomousControlLoop::new(config, snapshot, proposer, validator)
421    }
422
423    #[tokio::test]
424    async fn test_control_loop_creation() {
425        let loop_sys = create_test_loop();
426        assert_eq!(loop_sys.state().await, LoopState::Idle);
427    }
428
429    #[tokio::test]
430    async fn test_control_loop_observation() {
431        let loop_sys = create_test_loop();
432
433        let obs = Observation {
434            entity: "test_entity".to_string(),
435            properties: std::collections::BTreeMap::from([(
436                "type".to_string(),
437                "test".to_string(),
438            )]),
439            timestamp: 1000,
440            source: ObservationSource::Data,
441        };
442
443        loop_sys.observe(obs).await;
444
445        let miner = loop_sys.pattern_miner.read().await;
446        assert_eq!(miner.observations.len(), 1);
447    }
448
449    #[tokio::test]
450    async fn test_control_loop_iteration() {
451        let loop_sys = create_test_loop();
452
453        // Add some observations
454        for i in 0..3 {
455            let obs = Observation {
456                entity: format!("entity_{}", i),
457                properties: std::collections::BTreeMap::from([(
458                    "type".to_string(),
459                    "test".to_string(),
460                )]),
461                timestamp: 1000 + i as u64,
462                source: ObservationSource::Data,
463            };
464            loop_sys.observe(obs).await;
465        }
466
467        // Run one iteration
468        let telemetry = loop_sys.iterate().await.unwrap();
469        assert_eq!(telemetry.observation_count, 3);
470    }
471
472    #[tokio::test]
473    async fn test_control_loop_bounded_run() {
474        let loop_sys = create_test_loop();
475
476        // Add observations
477        for i in 0..5 {
478            let obs = Observation {
479                entity: format!("entity_{}", i),
480                properties: std::collections::BTreeMap::from([(
481                    "type".to_string(),
482                    "test".to_string(),
483                )]),
484                timestamp: 1000 + i as u64,
485                source: ObservationSource::Data,
486            };
487            loop_sys.observe(obs).await;
488        }
489
490        // Run 1 iteration
491        let result = loop_sys.run_bounded(1).await;
492        assert!(result.is_ok());
493
494        let telemetry = loop_sys.telemetry().await;
495        assert_eq!(telemetry.len(), 1);
496    }
497
498    #[tokio::test]
499    async fn test_control_loop_state_transitions() {
500        let loop_sys = create_test_loop();
501
502        // Add observations
503        for i in 0..3 {
504            let obs = Observation {
505                entity: format!("entity_{}", i),
506                properties: std::collections::BTreeMap::from([(
507                    "type".to_string(),
508                    "test".to_string(),
509                )]),
510                timestamp: 1000 + i as u64,
511                source: ObservationSource::Data,
512            };
513            loop_sys.observe(obs).await;
514        }
515
516        loop_sys.iterate().await.unwrap();
517
518        let final_state = loop_sys.state().await;
519        assert_eq!(final_state, LoopState::Recording);
520    }
521}