1use std::collections::HashSet;
23use std::sync::Arc;
24use std::time::Instant;
25
26use serde::{Deserialize, Serialize};
27use tokio_util::sync::CancellationToken;
28use zeph_common::SessionId;
29use zeph_common::timestamp;
30use zeph_llm::any::AnyProvider;
31use zeph_memory::semantic::SemanticMemory;
32use zeph_memory::store::experiments::NewExperimentResult;
33
34use super::error::EvalError;
35use super::evaluator::Evaluator;
36use super::generator::VariationGenerator;
37use super::snapshot::ConfigSnapshot;
38use super::types::{ExperimentResult, ExperimentSource, Variation};
39use zeph_config::ExperimentConfig;
40
41#[derive(Debug, Clone, Serialize, Deserialize)]
49pub struct ExperimentSessionReport {
50 pub session_id: SessionId,
52 pub results: Vec<ExperimentResult>,
54 pub best_config: ConfigSnapshot,
56 pub baseline_score: f64,
60 pub final_score: f64,
64 pub total_improvement: f64,
66 pub wall_time_ms: u64,
68 pub cancelled: bool,
73}
74
75pub struct ExperimentEngine {
96 evaluator: Evaluator,
97 generator: Box<dyn VariationGenerator>,
98 subject: Arc<AnyProvider>,
99 baseline: ConfigSnapshot,
100 config: ExperimentConfig,
101 memory: Option<Arc<SemanticMemory>>,
102 session_id: SessionId,
103 cancel: CancellationToken,
104 source: ExperimentSource,
105}
106
107const MAX_CONSECUTIVE_NAN: u32 = 3;
110
111impl ExperimentEngine {
112 pub fn new(
125 evaluator: Evaluator,
126 generator: Box<dyn VariationGenerator>,
127 subject: Arc<AnyProvider>,
128 baseline: ConfigSnapshot,
129 config: ExperimentConfig,
130 memory: Option<Arc<SemanticMemory>>,
131 ) -> Self {
132 Self {
133 evaluator,
134 generator,
135 subject,
136 baseline,
137 config,
138 memory,
139 session_id: SessionId::generate(),
140 cancel: CancellationToken::new(),
141 source: ExperimentSource::Manual,
142 }
143 }
144
145 #[must_use]
150 pub fn with_source(mut self, source: ExperimentSource) -> Self {
151 self.source = source;
152 self
153 }
154
155 #[must_use]
160 pub fn cancel_token(&self) -> CancellationToken {
161 self.cancel.clone()
162 }
163
164 pub fn stop(&self) {
168 self.cancel.cancel();
169 }
170
171 #[tracing::instrument(
192 name = "experiments.engine.run",
193 skip(self),
194 fields(session_id = %self.session_id, source = %self.source)
195 )]
196 pub async fn run(&mut self) -> Result<ExperimentSessionReport, EvalError> {
197 let start = Instant::now();
198 let best_snapshot = self.baseline.clone();
199
200 let baseline_report = tokio::select! {
203 biased;
204 () = self.cancel.cancelled() => {
205 tracing::info!(session_id = %self.session_id, "cancelled before baseline");
206 #[allow(clippy::cast_possible_truncation)]
207 return Ok(ExperimentSessionReport {
208 session_id: self.session_id.clone(),
209 results: vec![],
210 best_config: best_snapshot,
211 baseline_score: f64::NAN,
212 final_score: f64::NAN,
213 total_improvement: 0.0,
214 wall_time_ms: start.elapsed().as_millis() as u64,
215 cancelled: true,
216 });
217 }
218 report = self.evaluator.evaluate(&self.subject) => report?,
219 };
220
221 let initial_baseline_score = baseline_report.mean_score;
223 if initial_baseline_score.is_nan() {
224 return Err(EvalError::Storage(
225 "baseline evaluation produced NaN mean score; \
226 check evaluator budget and judge responses"
227 .into(),
228 ));
229 }
230 tracing::info!(
231 session_id = %self.session_id,
232 baseline_score = initial_baseline_score,
233 "experiment session started"
234 );
235 self.run_loop(start, initial_baseline_score, best_snapshot)
236 .await
237 }
238
239 #[allow(clippy::too_many_lines)] #[tracing::instrument(
246 name = "experiments.engine.run_loop",
247 skip(self, start, best_snapshot),
248 fields(session_id = %self.session_id, source = %self.source)
249 )]
250 async fn run_loop(
251 &mut self,
252 start: Instant,
253 initial_baseline_score: f64,
254 mut best_snapshot: ConfigSnapshot,
255 ) -> Result<ExperimentSessionReport, EvalError> {
256 let wall_limit = std::time::Duration::from_secs(self.config.max_wall_time_secs);
257 let mut results: Vec<ExperimentResult> = Vec::new();
258 let mut visited: HashSet<Variation> = HashSet::new();
259 let (mut best_score, mut consecutive_nan) = (initial_baseline_score, 0u32);
260
261 'main: loop {
262 if results.len() >= self.config.max_experiments as usize {
263 tracing::info!(session_id = %self.session_id, "budget exhausted");
264 break;
265 }
266 if start.elapsed() >= wall_limit {
267 tracing::info!(session_id = %self.session_id, "wall-time limit reached");
268 break;
269 }
270 let Some(variation) = self.generator.next(&best_snapshot, &visited) else {
271 tracing::info!(session_id = %self.session_id, "search space exhausted");
272 break;
273 };
274 visited.insert(variation.clone());
275 let candidate_snapshot = best_snapshot.apply(&variation);
276 let patched = (*self.subject)
277 .clone()
278 .with_generation_overrides(candidate_snapshot.to_generation_overrides());
279 let candidate_report = tokio::select! {
280 biased;
281 () = self.cancel.cancelled() => {
282 tracing::info!(session_id = %self.session_id, "experiment cancelled");
283 break 'main;
284 }
285 report = self.evaluator.evaluate(&patched) => report?,
286 };
287 if candidate_report.mean_score.is_nan() {
288 consecutive_nan += 1;
289 tracing::warn!(
290 session_id = %self.session_id, param = %variation.parameter,
291 is_partial = candidate_report.is_partial, consecutive_nan,
292 "NaN mean score — skipping variation"
293 );
294 if consecutive_nan >= MAX_CONSECUTIVE_NAN {
295 tracing::warn!(session_id = %self.session_id, "consecutive NaN cap reached");
296 break;
297 }
298 continue;
299 }
300 consecutive_nan = 0;
301 let candidate_score = candidate_report.mean_score;
302 let delta = candidate_score - best_score;
303 let accepted = delta >= self.config.min_improvement;
304 let result_id = self
305 .persist_result(
306 &variation,
307 best_score,
308 candidate_score,
309 delta,
310 accepted,
311 candidate_report.p50_latency_ms,
312 candidate_report.total_tokens,
313 )
314 .await?;
315 let pre_accept_baseline = best_score;
316 self.log_outcome(&variation, delta, accepted, best_score);
317 if accepted {
318 best_snapshot = candidate_snapshot;
319 best_score = candidate_score;
320 }
321 results.push(ExperimentResult {
322 id: result_id,
323 session_id: self.session_id.clone(),
324 variation,
325 baseline_score: pre_accept_baseline,
326 candidate_score,
327 delta,
328 latency_ms: candidate_report.p50_latency_ms,
329 tokens_used: candidate_report.total_tokens,
330 accepted,
331 source: self.source.clone(),
332 created_at: timestamp::utc_now_rfc3339(),
333 });
334 }
335
336 #[allow(clippy::cast_possible_truncation)]
337 let wall_time_ms = start.elapsed().as_millis() as u64;
338 let total_improvement = best_score - initial_baseline_score;
339 tracing::info!(
340 session_id = %self.session_id, total = results.len(),
341 baseline_score = initial_baseline_score, final_score = best_score,
342 total_improvement, wall_time_ms, cancelled = self.cancel.is_cancelled(),
343 "experiment session complete"
344 );
345 Ok(ExperimentSessionReport {
346 session_id: self.session_id.clone(),
347 results,
348 best_config: best_snapshot,
349 baseline_score: initial_baseline_score,
350 final_score: best_score,
351 total_improvement,
352 wall_time_ms,
353 cancelled: self.cancel.is_cancelled(),
354 })
355 }
356
357 #[tracing::instrument(name = "experiments.engine.persist_result", skip_all)]
366 #[allow(clippy::too_many_arguments)] async fn persist_result(
368 &self,
369 variation: &Variation,
370 baseline_score: f64,
371 candidate_score: f64,
372 delta: f64,
373 accepted: bool,
374 p50_latency_ms: u64,
375 total_tokens: u64,
376 ) -> Result<Option<i64>, EvalError> {
377 let Some(mem) = &self.memory else {
378 return Ok(None);
379 };
380 let value_json = serde_json::to_string(&variation.value)
381 .map_err(|e| EvalError::Storage(e.to_string()))?;
382 #[allow(clippy::cast_possible_wrap)]
383 let new_result = NewExperimentResult {
384 session_id: self.session_id.as_str(),
385 parameter: variation.parameter.as_str(),
386 value_json: &value_json,
387 baseline_score,
388 candidate_score,
389 delta,
390 latency_ms: p50_latency_ms as i64,
391 tokens_used: total_tokens as i64,
392 accepted,
393 source: self.source.as_str(),
394 };
395 mem.sqlite()
396 .insert_experiment_result(&new_result)
397 .await
398 .map(Some)
399 .map_err(|e: zeph_memory::error::MemoryError| EvalError::Storage(e.to_string()))
400 }
401
402 fn log_outcome(&self, variation: &Variation, delta: f64, accepted: bool, new_score: f64) {
403 if accepted {
404 tracing::info!(
405 session_id = %self.session_id,
406 param = %variation.parameter,
407 value = %variation.value,
408 delta,
409 new_best_score = new_score,
410 "variation accepted — new baseline"
411 );
412 } else {
413 tracing::info!(
414 session_id = %self.session_id,
415 param = %variation.parameter,
416 value = %variation.value,
417 delta,
418 "variation rejected"
419 );
420 }
421 }
422}
423
424#[cfg(test)]
425mod tests {
426 #![allow(clippy::doc_markdown)]
427
428 use super::*;
429 use crate::benchmark::{BenchmarkCase, BenchmarkSet};
430 use crate::evaluator::Evaluator;
431 use crate::generator::VariationGenerator;
432 use crate::snapshot::ConfigSnapshot;
433 use crate::types::{ParameterKind, Variation, VariationValue};
434 use ordered_float::OrderedFloat;
435 use std::sync::Arc;
436 use zeph_config::ExperimentConfig;
437
438 fn make_benchmark() -> BenchmarkSet {
439 BenchmarkSet {
440 cases: vec![BenchmarkCase {
441 prompt: "What is 2+2?".into(),
442 context: None,
443 reference: None,
444 tags: None,
445 }],
446 }
447 }
448
449 fn default_config() -> ExperimentConfig {
450 ExperimentConfig {
451 max_experiments: 10,
452 max_wall_time_secs: 3600,
453 min_improvement: 0.0,
454 ..Default::default()
455 }
456 }
457
458 struct NVariationGenerator {
460 variations: Vec<Variation>,
461 pos: usize,
462 }
463
464 impl NVariationGenerator {
465 fn new(n: usize) -> Self {
466 let variations = (0..n)
467 .map(|i| Variation {
468 parameter: ParameterKind::Temperature,
469 #[allow(clippy::cast_precision_loss)]
470 value: VariationValue::Float(OrderedFloat(0.5 + i as f64 * 0.1)),
471 })
472 .collect();
473 Self { variations, pos: 0 }
474 }
475 }
476
477 impl VariationGenerator for NVariationGenerator {
478 fn next(
479 &mut self,
480 _baseline: &ConfigSnapshot,
481 visited: &HashSet<Variation>,
482 ) -> Option<Variation> {
483 while self.pos < self.variations.len() {
484 let v = self.variations[self.pos].clone();
485 self.pos += 1;
486 if !visited.contains(&v) {
487 return Some(v);
488 }
489 }
490 None
491 }
492
493 fn name(&self) -> &'static str {
494 "n_variation"
495 }
496 }
497
498 #[cfg(test)]
499 fn make_subject_mock(n_responses: usize) -> zeph_llm::any::AnyProvider {
500 use zeph_llm::any::AnyProvider;
501 use zeph_llm::mock::MockProvider;
502
503 let responses: Vec<String> = (0..n_responses).map(|_| "Four".to_string()).collect();
507 AnyProvider::Mock(MockProvider::with_responses(responses))
508 }
509
510 #[cfg(test)]
511 fn make_judge_mock(n_responses: usize) -> zeph_llm::any::AnyProvider {
512 use zeph_llm::any::AnyProvider;
513 use zeph_llm::mock::MockProvider;
514
515 let responses: Vec<String> = (0..n_responses)
516 .map(|_| r#"{"score": 8.0, "reason": "correct"}"#.to_string())
517 .collect();
518 AnyProvider::Mock(MockProvider::with_responses(responses))
519 }
520
521 #[cfg(test)]
522 #[tokio::test]
523 async fn engine_completes_with_no_accepted_variations() {
524 let config = ExperimentConfig {
526 max_experiments: 10,
527 max_wall_time_secs: 3600,
528 min_improvement: 100.0,
529 ..Default::default()
530 };
531 let subject = make_subject_mock(2);
533 let judge = make_judge_mock(2);
534 let evaluator = Evaluator::new(Arc::new(judge), make_benchmark(), 1_000_000).unwrap();
535
536 let mut engine = ExperimentEngine::new(
537 evaluator,
538 Box::new(NVariationGenerator::new(1)),
539 Arc::new(subject),
540 ConfigSnapshot::default(),
541 config,
542 None,
543 );
544
545 let report = engine.run().await.unwrap();
546 assert_eq!(report.results.len(), 1);
547 assert!(!report.results[0].accepted);
548 assert!(!report.session_id.is_empty());
549 assert!(!report.cancelled);
550 }
551
552 #[cfg(test)]
553 #[tokio::test]
554 async fn engine_respects_max_experiments() {
555 let config = ExperimentConfig {
556 max_experiments: 3,
557 max_wall_time_secs: 3600,
558 min_improvement: 0.0,
559 ..Default::default()
560 };
561 let subject = make_subject_mock(4);
564 let judge = make_judge_mock(4);
565 let evaluator = Evaluator::new(Arc::new(judge), make_benchmark(), 1_000_000).unwrap();
566
567 let mut engine = ExperimentEngine::new(
568 evaluator,
569 Box::new(NVariationGenerator::new(5)),
570 Arc::new(subject),
571 ConfigSnapshot::default(),
572 config,
573 None,
574 );
575
576 let report = engine.run().await.unwrap();
577 assert_eq!(report.results.len(), 3);
578 assert!(!report.cancelled);
579 }
580
581 #[cfg(test)]
582 #[tokio::test]
583 async fn engine_cancellation_before_baseline() {
584 let config = ExperimentConfig {
586 max_experiments: 100,
587 max_wall_time_secs: 3600,
588 min_improvement: 0.0,
589 ..Default::default()
590 };
591 let subject = make_subject_mock(2);
592 let judge = make_judge_mock(2);
593 let evaluator = Evaluator::new(Arc::new(judge), make_benchmark(), 1_000_000).unwrap();
594
595 let mut engine = ExperimentEngine::new(
596 evaluator,
597 Box::new(NVariationGenerator::new(100)),
598 Arc::new(subject),
599 ConfigSnapshot::default(),
600 config,
601 None,
602 );
603 engine.stop(); let report = engine.run().await.unwrap();
605 assert!(report.cancelled);
606 assert!(report.results.is_empty());
607 }
608
609 #[cfg(test)]
610 #[tokio::test]
611 async fn engine_cancellation_stops_loop() {
612 let config = ExperimentConfig {
621 max_experiments: 10,
622 max_wall_time_secs: 3600,
623 min_improvement: 0.0,
624 ..Default::default()
625 };
626 let subject = make_subject_mock(2);
627 let judge = make_judge_mock(2);
628 let evaluator = Evaluator::new(Arc::new(judge), make_benchmark(), 1_000_000).unwrap();
629
630 let mut engine = ExperimentEngine::new(
631 evaluator,
632 Box::new(NVariationGenerator::new(10)),
633 Arc::new(subject),
634 ConfigSnapshot::default(),
635 config,
636 None,
637 );
638
639 let external_token = engine.cancel_token();
641 assert!(!external_token.is_cancelled());
642 engine.stop();
643 assert!(
644 external_token.is_cancelled(),
645 "cancel_token() must share the same token"
646 );
647
648 let report = engine.run().await.unwrap();
649 assert!(report.cancelled);
650 }
651
652 #[cfg(test)]
653 #[tokio::test]
654 async fn engine_progressive_baseline_updates() {
655 let config = ExperimentConfig {
658 max_experiments: 1,
659 max_wall_time_secs: 3600,
660 min_improvement: 0.0,
661 ..Default::default()
662 };
663 let subject = make_subject_mock(2);
665 let judge = make_judge_mock(2);
666 let evaluator = Evaluator::new(Arc::new(judge), make_benchmark(), 1_000_000).unwrap();
667
668 let initial_baseline = ConfigSnapshot::default();
669 let mut engine = ExperimentEngine::new(
670 evaluator,
671 Box::new(NVariationGenerator::new(1)),
672 Arc::new(subject),
673 initial_baseline.clone(),
674 config,
675 None,
676 );
677
678 let report = engine.run().await.unwrap();
679 assert_eq!(report.results.len(), 1);
680 assert!(report.results[0].accepted, "variation should be accepted");
681 assert!(
683 (report.best_config.temperature - initial_baseline.temperature).abs() > 1e-9,
684 "best_config.temperature should have changed after accepted variation"
685 );
686 assert!(!report.baseline_score.is_nan());
687 assert!(!report.final_score.is_nan());
688 assert!(
690 (report.results[0].baseline_score - report.baseline_score).abs() < 1e-9,
691 "result.baseline_score must equal initial baseline_score (pre-acceptance)"
692 );
693 }
694
695 #[cfg(test)]
696 #[tokio::test]
697 async fn engine_handles_search_space_exhaustion() {
698 let config = default_config();
699 let subject = make_subject_mock(1);
702 let judge = make_judge_mock(1);
703 let evaluator = Evaluator::new(Arc::new(judge), make_benchmark(), 1_000_000).unwrap();
704
705 let mut engine = ExperimentEngine::new(
706 evaluator,
707 Box::new(NVariationGenerator::new(0)),
708 Arc::new(subject),
709 ConfigSnapshot::default(),
710 config,
711 None,
712 );
713
714 let report = engine.run().await.unwrap();
715 assert!(report.results.is_empty());
716 assert!(!report.cancelled);
717 }
718
719 #[cfg(test)]
720 #[tokio::test]
721 async fn engine_skips_nan_scores() {
722 use zeph_llm::any::AnyProvider;
723 use zeph_llm::mock::MockProvider;
724
725 let config = ExperimentConfig {
731 max_experiments: 5,
732 max_wall_time_secs: 3600,
733 min_improvement: 0.0,
734 ..Default::default()
735 };
736 let subject = AnyProvider::Mock(MockProvider::with_responses(vec![
738 "A".into(),
739 "A".into(),
740 "A".into(),
741 "A".into(),
742 ]));
743 let judge = AnyProvider::Mock(MockProvider::with_responses(vec![
745 r#"{"score": 8.0, "reason": "ok"}"#.into(),
746 ]));
747 let evaluator = Evaluator::new(Arc::new(judge), make_benchmark(), 1_000_000).unwrap();
749
750 let mut engine = ExperimentEngine::new(
751 evaluator,
752 Box::new(NVariationGenerator::new(5)),
753 Arc::new(subject),
754 ConfigSnapshot::default(),
755 config,
756 None,
757 );
758
759 let report = engine.run().await.unwrap();
761 assert!(
763 report.results.is_empty(),
764 "all NaN iterations should be skipped"
765 );
766 assert!(!report.cancelled);
767 }
768
769 #[cfg(test)]
770 #[tokio::test]
771 async fn engine_nan_baseline_returns_error() {
772 use zeph_llm::any::AnyProvider;
773 use zeph_llm::mock::MockProvider;
774
775 let config = ExperimentConfig {
777 max_experiments: 5,
778 max_wall_time_secs: 3600,
779 min_improvement: 0.0,
780 ..Default::default()
781 };
782 let subject = AnyProvider::Mock(MockProvider::with_responses(vec!["A".into()]));
784 let judge = AnyProvider::Mock(MockProvider::with_responses(vec![]));
786 let evaluator = Evaluator::new(Arc::new(judge), make_benchmark(), 1_000_000).unwrap();
787
788 let mut engine = ExperimentEngine::new(
789 evaluator,
790 Box::new(NVariationGenerator::new(5)),
791 Arc::new(subject),
792 ConfigSnapshot::default(),
793 config,
794 None,
795 );
796
797 let result = engine.run().await;
798 assert!(result.is_err(), "NaN baseline should return an error");
799 let err = result.unwrap_err();
800 assert!(
801 matches!(err, EvalError::Storage(_)),
802 "expected EvalError::Storage, got: {err:?}"
803 );
804 }
805
806 #[cfg(test)]
807 #[tokio::test]
808 async fn engine_persists_results_to_sqlite() {
809 use zeph_memory::testing::mock_semantic_memory;
810
811 let memory = mock_semantic_memory().await.unwrap();
812 let config = ExperimentConfig {
813 max_experiments: 1,
814 max_wall_time_secs: 3600,
815 min_improvement: 0.0,
816 ..Default::default()
817 };
818 let subject = make_subject_mock(2);
820 let judge = make_judge_mock(2);
821 let evaluator = Evaluator::new(Arc::new(judge), make_benchmark(), 1_000_000).unwrap();
822
823 let session_id = {
824 let mut engine = ExperimentEngine::new(
825 evaluator,
826 Box::new(NVariationGenerator::new(1)),
827 Arc::new(subject),
828 ConfigSnapshot::default(),
829 config,
830 Some(Arc::clone(&memory)),
831 );
832 engine.run().await.unwrap();
833 engine.session_id.clone()
834 };
835
836 let rows = memory
837 .sqlite()
838 .list_experiment_results(Some(&session_id), 10)
839 .await
840 .unwrap();
841 assert_eq!(rows.len(), 1, "expected one persisted result");
842 assert_eq!(rows[0].session_id, session_id.as_str());
843 }
844
845 #[test]
846 fn session_report_serde_roundtrip() {
847 let report = ExperimentSessionReport {
848 session_id: SessionId::new("test-session"),
849 results: vec![],
850 best_config: ConfigSnapshot::default(),
851 baseline_score: 7.5,
852 final_score: 8.0,
853 total_improvement: 0.5,
854 wall_time_ms: 1_234,
855 cancelled: false,
856 };
857 let json = serde_json::to_string(&report).expect("serialize");
858 let report2: ExperimentSessionReport = serde_json::from_str(&json).expect("deserialize");
859 assert_eq!(report2.session_id, report.session_id);
860 assert!((report2.baseline_score - report.baseline_score).abs() < f64::EPSILON);
861 assert!((report2.final_score - report.final_score).abs() < f64::EPSILON);
862 assert_eq!(report2.wall_time_ms, report.wall_time_ms);
863 assert!(!report2.cancelled);
864 }
865
866 #[test]
867 fn utc_now_rfc3339_format() {
868 let s = timestamp::utc_now_rfc3339();
869 assert_eq!(s.len(), 20, "timestamp must be 20 chars (RFC 3339): {s}");
870 assert_eq!(&s[4..5], "-");
871 assert_eq!(&s[7..8], "-");
872 assert_eq!(&s[10..11], "T");
873 assert_eq!(&s[13..14], ":");
874 assert_eq!(&s[16..17], ":");
875 assert!(s.ends_with('Z'));
876 }
877
878 #[test]
880 fn utc_now_rfc3339_is_non_empty() {
881 let ts = timestamp::utc_now_rfc3339();
882 assert!(!ts.is_empty());
883 assert_eq!(ts.len(), 20);
884 }
885
886 #[tokio::test]
887 async fn experiment_result_created_at_is_rfc3339() {
888 let config = ExperimentConfig {
889 max_experiments: 1,
890 max_wall_time_secs: 3600,
891 min_improvement: 0.0,
892 ..Default::default()
893 };
894 let subject = make_subject_mock(2);
895 let judge = make_judge_mock(2);
896 let evaluator = Evaluator::new(Arc::new(judge), make_benchmark(), 1_000_000).unwrap();
897
898 let mut engine = ExperimentEngine::new(
899 evaluator,
900 Box::new(NVariationGenerator::new(1)),
901 Arc::new(subject),
902 ConfigSnapshot::default(),
903 config,
904 None,
905 );
906
907 let report = engine.run().await.unwrap();
908 assert_eq!(report.results.len(), 1);
909 let created_at = &report.results[0].created_at;
910 assert!(!created_at.is_empty(), "created_at must not be empty");
911 assert_eq!(
912 created_at.len(),
913 20,
914 "RFC 3339 timestamp must be 20 chars: {created_at}"
915 );
916 assert!(
917 created_at.contains('T'),
918 "RFC 3339 timestamp must contain 'T': {created_at}"
919 );
920 assert!(
921 created_at.ends_with('Z'),
922 "RFC 3339 timestamp must end with 'Z': {created_at}"
923 );
924 }
925
926 #[test]
928 fn experiment_engine_is_send() {
929 fn assert_send<T: Send>() {}
930 let _ = assert_send::<ExperimentEngine>;
933 }
934
935 #[tokio::test]
936 async fn engine_with_source_scheduled_propagates_to_results() {
937 let config = ExperimentConfig {
938 max_experiments: 1,
939 max_wall_time_secs: 3600,
940 min_improvement: 0.0,
941 ..Default::default()
942 };
943 let subject = make_subject_mock(2);
944 let judge = make_judge_mock(2);
945 let evaluator = Evaluator::new(Arc::new(judge), make_benchmark(), 1_000_000).unwrap();
946
947 let mut engine = ExperimentEngine::new(
948 evaluator,
949 Box::new(NVariationGenerator::new(1)),
950 Arc::new(subject),
951 ConfigSnapshot::default(),
952 config,
953 None,
954 )
955 .with_source(ExperimentSource::Scheduled);
956
957 let report = engine.run().await.unwrap();
958 assert_eq!(report.results.len(), 1);
959 assert_eq!(
960 report.results[0].source,
961 ExperimentSource::Scheduled,
962 "with_source(Scheduled) must propagate to ExperimentResult"
963 );
964 }
965}