1use serde::{Deserialize, Serialize};
2use 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#[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#[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#[derive(Debug, Clone)]
48pub struct ControlLoopConfig {
49 pub iteration_interval_ms: u64,
51
52 pub max_iterations: Option<usize>,
54
55 pub auto_promote: bool,
57
58 pub sector: String,
60
61 pub min_proposal_confidence: f64,
63
64 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
81pub 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 pub async fn state(&self) -> LoopState {
115 *self.state.read().await
116 }
117
118 pub async fn telemetry(&self) -> Vec<IterationTelemetry> {
120 self.telemetry.read().await.clone()
121 }
122
123 pub async fn observe(&self, obs: Observation) {
125 let mut miner = self.pattern_miner.write().await;
126 miner.add_observation(obs);
127 }
128
129 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 let miner = self.pattern_miner.read().await;
145 telemetry.observation_count = miner.observation_count();
146 drop(miner);
147
148 *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 *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 let valid_proposals: Vec<_> = proposals
178 .iter()
179 .filter(|p| p.confidence >= self.config.min_proposal_confidence)
180 .cloned()
181 .collect();
182
183 *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 let mut new_triples = current_snap.snapshot().triples.as_ref().clone();
193
194 for triple_pattern in &proposal.triples_to_remove {
196 new_triples.retain(|stmt| !stmt.subject.contains(triple_pattern));
197 }
198
199 for triple_str in &proposal.triples_to_add {
201 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 if static_ev.passed && dynamic_ev.passed && perf_ev.passed {
243 *self.state.write().await = LoopState::Promoting;
245
246 if self.config.auto_promote {
247 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 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 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 *self.state.write().await = LoopState::Recording;
306 telemetry.total_duration_ms = get_time_ms() - start_ms;
307
308 Ok(telemetry)
309 }
310
311 pub async fn run(&self) -> Result<(), String> {
313 *self.state.write().await = LoopState::Idle;
314
315 let mut iteration = 0;
316 loop {
317 if let Some(max) = self.config.max_iterations {
319 if iteration >= max {
320 break;
321 }
322 }
323
324 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 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 pub async fn run_bounded(&self, max_iters: usize) -> Result<(), String> {
347 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 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
376fn 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 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 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 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 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 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}