Skip to main content

everruns_core/
in_memory_loop.rs

1// In-Memory Agentic Loop
2//
3// Convenience helpers for running full agentic loops in memory without
4// external dependencies (database, real LLM, etc.). Perfect for:
5// - Unit and integration tests
6// - Prototyping and experimentation
7// - Examples and documentation
8//
9// The `InMemoryAgenticLoop` bundles all in-memory stores and atoms,
10// providing a simple API for executing agent turns.
11
12use std::sync::Arc;
13
14use async_trait::async_trait;
15use chrono::Utc;
16
17use crate::agent::{Agent, AgentStatus};
18use crate::atoms::{
19    ActAtom, ActInput, Atom, AtomContext, InputAtom, InputAtomInput, ReasonAtom, ReasonInput,
20};
21use crate::capabilities::{AgentCapabilityConfig, Capability, CapabilityRegistry};
22use crate::driver_registry::{DriverId, DriverRegistry};
23use crate::error::Result;
24use crate::events::{Event, EventData, EventRequest, OUTPUT_MESSAGE_COMPLETED};
25use crate::in_memory::{
26    InMemoryAgentStore, InMemoryEventEmitter, InMemoryHarnessStore, InMemoryMessageRetriever,
27    InMemoryProviderStore, InMemorySessionStore,
28};
29use crate::llmsim_driver::{LlmSimConfig, LlmSimDriver};
30use crate::message::Message;
31use crate::message_retriever::{InputMessage, MessageRetriever};
32use crate::session::{Session, SessionStatus};
33use crate::tool_types::ToolCall;
34use crate::tools::{Tool, ToolRegistry, ToolRegistryBuilder};
35use crate::traits::{EventEmitter, ResolvedModel};
36use crate::turn::{TurnAction, TurnContext, TurnOutcome, TurnStateMachine};
37use crate::typed_id::{AgentId, HarnessId, SessionId, TurnId};
38
39// ============================================================================
40// Bridging Event Emitter
41// ============================================================================
42
43/// Event emitter that bridges events to message storage
44///
45/// When a `output.message.completed` event is emitted, it also stores the message
46/// in the provided message retriever. This enables full agentic loops
47/// in memory without the database layer.
48#[derive(Clone)]
49struct BridgingEventEmitter {
50    inner: InMemoryEventEmitter,
51    message_retriever: InMemoryMessageRetriever,
52}
53
54impl BridgingEventEmitter {
55    fn new(message_retriever: InMemoryMessageRetriever) -> Self {
56        Self {
57            inner: InMemoryEventEmitter::new(),
58            message_retriever,
59        }
60    }
61
62    async fn events(&self) -> Vec<Event> {
63        self.inner.events().await
64    }
65
66    async fn events_by_type(&self, event_type: &str) -> Vec<Event> {
67        self.inner.events_by_type(event_type).await
68    }
69
70    async fn event_count(&self) -> usize {
71        self.inner.event_count().await
72    }
73
74    async fn clear(&self) {
75        self.inner.clear().await;
76    }
77}
78
79#[async_trait]
80impl EventEmitter for BridgingEventEmitter {
81    async fn emit(&self, request: EventRequest) -> Result<Event> {
82        // If this is an output.message.completed event, also store the message
83        if request.data.event_type() == OUTPUT_MESSAGE_COMPLETED
84            && let EventData::OutputMessageCompleted(data) = &request.data
85        {
86            // Store the message in the retriever
87            let _ = self
88                .message_retriever
89                .store(request.session_id, data.message.clone())
90                .await;
91        }
92
93        // Delegate to the inner emitter
94        self.inner.emit(request).await
95    }
96}
97
98// ============================================================================
99// Turn Result
100// ============================================================================
101
102/// Result of executing a turn
103#[derive(Debug, Clone)]
104pub struct TurnResult {
105    /// Final text response from the agent
106    pub response: String,
107    /// Number of reasoning iterations (Reason → Act cycles)
108    pub iterations: usize,
109    /// Total tool calls made during the turn
110    pub tool_calls_count: usize,
111    /// Whether the turn completed successfully
112    pub success: bool,
113    /// Error message if the turn failed
114    pub error: Option<String>,
115    /// Turn ID for this turn
116    pub turn_id: TurnId,
117}
118
119impl TurnResult {
120    /// Check if the response contains a specific substring
121    pub fn contains(&self, text: &str) -> bool {
122        self.response.contains(text)
123    }
124
125    /// Create a TurnResult from a TurnOutcome and turn_id.
126    fn from_outcome(outcome: TurnOutcome, turn_id: TurnId) -> Self {
127        match outcome {
128            TurnOutcome::Success {
129                response,
130                iterations,
131                tool_calls_count,
132            } => Self {
133                response,
134                iterations,
135                tool_calls_count,
136                success: true,
137                error: None,
138                turn_id,
139            },
140            TurnOutcome::Failed { error, iterations } => Self {
141                response: String::new(),
142                iterations,
143                tool_calls_count: 0,
144                success: false,
145                error: Some(error),
146                turn_id,
147            },
148            TurnOutcome::MaxIterationsReached {
149                response,
150                iterations,
151                tool_calls_count,
152            } => Self {
153                response,
154                iterations,
155                tool_calls_count,
156                success: true, // Max iterations is not a failure
157                error: None,
158                turn_id,
159            },
160            // A sealed turn was deliberately stopped (EVE-534). Surface it as a
161            // non-success with the seal reason so in-memory callers can observe
162            // it distinctly from a normal completion.
163            TurnOutcome::Sealed {
164                reason,
165                response,
166                iterations,
167                tool_calls_count,
168            } => Self {
169                response,
170                iterations,
171                tool_calls_count,
172                success: false,
173                error: Some(format!("turn sealed: {reason}")),
174                turn_id,
175            },
176        }
177    }
178}
179
180// ============================================================================
181// Builder
182// ============================================================================
183
184/// Builder for creating an `InMemoryAgenticLoop`
185pub struct InMemoryAgenticLoopBuilder {
186    agent_name: String,
187    system_prompt: String,
188    model: Option<ResolvedModel>,
189    driver_registry: Option<DriverRegistry>,
190    llm_sim_config: Option<LlmSimConfig>,
191    tools: Vec<Box<dyn Tool>>,
192    capabilities: Vec<Box<dyn Capability>>,
193    max_iterations: usize,
194    parallel_tool_calls: Option<bool>,
195    reasoning_effort_handle: Option<crate::traits::ReasoningEffortHandle>,
196}
197
198impl Default for InMemoryAgenticLoopBuilder {
199    fn default() -> Self {
200        Self::new()
201    }
202}
203
204impl InMemoryAgenticLoopBuilder {
205    /// Create a new builder with defaults (uses simulated LLM)
206    pub fn new() -> Self {
207        Self {
208            agent_name: "Test Agent".to_string(),
209            system_prompt: "You are a helpful assistant.".to_string(),
210            model: None,
211            driver_registry: None,
212            llm_sim_config: Some(LlmSimConfig::default()),
213            tools: vec![],
214            capabilities: vec![],
215            max_iterations: 10,
216            parallel_tool_calls: None,
217            reasoning_effort_handle: None,
218        }
219    }
220
221    /// Share a live reasoning-effort handle (EVE-595) across the loop's
222    /// `ReasonAtom` and `ActAtom`. Tools receive a clone via their
223    /// `ToolContext` and can mutate it mid-turn so subsequent LLM steps in the
224    /// same `run_turn` observe the new effort.
225    pub fn reasoning_effort_handle(mut self, handle: crate::traits::ReasoningEffortHandle) -> Self {
226        self.reasoning_effort_handle = Some(handle);
227        self
228    }
229
230    /// Set the agent name
231    pub fn agent_name(mut self, name: impl Into<String>) -> Self {
232        self.agent_name = name.into();
233        self
234    }
235
236    /// Set the system prompt
237    pub fn system_prompt(mut self, prompt: impl Into<String>) -> Self {
238        self.system_prompt = prompt.into();
239        self
240    }
241
242    /// Use a simulated LLM with a fixed response (no real API calls)
243    pub fn with_simulated_response(mut self, response: impl Into<String>) -> Self {
244        self.llm_sim_config = Some(LlmSimConfig::fixed(response));
245        self.model = None;
246        self.driver_registry = None;
247        self
248    }
249
250    /// Use a simulated LLM with custom configuration
251    pub fn with_llm_sim(mut self, config: LlmSimConfig) -> Self {
252        self.llm_sim_config = Some(config);
253        self.model = None;
254        self.driver_registry = None;
255        self
256    }
257
258    /// Set the LLM model to use
259    ///
260    /// # Example
261    ///
262    /// ```ignore
263    /// use everruns_core::traits::ResolvedModel;
264    /// use everruns_core::provider::DriverId;
265    ///
266    /// let model = ResolvedModel {
267    ///     model: "claude-sonnet-4-20250514".to_string(),
268    ///     provider_type: DriverId::Anthropic,
269    ///     api_key: Some(std::env::var("ANTHROPIC_API_KEY").unwrap()),
270    ///     base_url: None,
271    /// };
272    ///
273    /// let runner = InMemoryAgenticLoop::builder()
274    ///     .model(model)
275    ///     .driver_registry(driver_registry)
276    ///     .build()
277    ///     .await?;
278    /// ```
279    pub fn model(mut self, model: ResolvedModel) -> Self {
280        self.model = Some(model);
281        self.llm_sim_config = None;
282        self
283    }
284
285    /// Set the driver registry for LLM providers
286    ///
287    /// # Example
288    ///
289    /// ```ignore
290    /// use everruns_core::driver_registry::DriverRegistry;
291    ///
292    /// let mut driver_registry = DriverRegistry::new();
293    /// everruns_anthropic::register_driver(&mut driver_registry);
294    ///
295    /// let runner = InMemoryAgenticLoop::builder()
296    ///     .model(model)
297    ///     .driver_registry(driver_registry)
298    ///     .build()
299    ///     .await?;
300    /// ```
301    pub fn driver_registry(mut self, driver_registry: DriverRegistry) -> Self {
302        self.driver_registry = Some(driver_registry);
303        self.llm_sim_config = None;
304        self
305    }
306
307    /// Add a tool
308    pub fn tool<T: Tool + 'static>(mut self, tool: T) -> Self {
309        self.tools.push(Box::new(tool));
310        self
311    }
312
313    /// Add a capability (which may provide tools and system prompt additions)
314    ///
315    /// Capabilities provide a way to bundle related tools and functionality.
316    /// For example, the `current_time` capability provides a `get_current_time` tool.
317    ///
318    /// # Example
319    ///
320    /// ```ignore
321    /// use everruns_core::capabilities::current_time::CurrentTimeCapability;
322    ///
323    /// let runner = InMemoryAgenticLoop::builder()
324    ///     .capability(CurrentTimeCapability)
325    ///     .build()
326    ///     .await?;
327    /// ```
328    pub fn capability<C: Capability + 'static>(mut self, capability: C) -> Self {
329        self.capabilities.push(Box::new(capability));
330        self
331    }
332
333    /// Set maximum iterations per turn
334    pub fn max_iterations(mut self, max: usize) -> Self {
335        self.max_iterations = max;
336        self
337    }
338
339    /// Set the request-level parallel tool calling preference (EVE-598).
340    ///
341    /// `Some(true)` signals the provider that parallel tool calls are wanted;
342    /// `Some(false)` requests at most one tool call per turn and forces serial
343    /// execution. `None` (default) preserves provider defaults.
344    pub fn parallel_tool_calls(mut self, parallel_tool_calls: Option<bool>) -> Self {
345        self.parallel_tool_calls = parallel_tool_calls;
346        self
347    }
348
349    /// Build the agentic loop
350    pub async fn build(self) -> Result<InMemoryAgenticLoop> {
351        // Create stores
352        let harness_store = InMemoryHarnessStore::new();
353        let agent_store = InMemoryAgentStore::new();
354        let session_store = InMemorySessionStore::new();
355        let message_retriever = InMemoryMessageRetriever::new();
356        let event_emitter = BridgingEventEmitter::new(message_retriever.clone());
357
358        // Build capability configs for the agent from capabilities
359        let agent_capability_configs: Vec<AgentCapabilityConfig> = self
360            .capabilities
361            .iter()
362            .map(|cap| AgentCapabilityConfig::new(cap.id()))
363            .collect();
364
365        // Create harness
366        let harness_id = HarnessId::new();
367        let now = Utc::now();
368        let harness = crate::harness::Harness {
369            id: harness_id,
370            name: "in-memory".to_string(),
371            display_name: Some("In-Memory Harness".to_string()),
372            description: None,
373            system_prompt: Some(self.system_prompt.clone()),
374            parent_harness_id: None,
375            default_model_id: None,
376            tags: vec![],
377            capabilities: vec![],
378            mcp_servers: Default::default(),
379            initial_files: vec![],
380            network_access: None,
381            parallel_tool_calls: None,
382            embedder_metadata: Default::default(),
383            is_built_in: false,
384            status: crate::harness::HarnessStatus::Active,
385            created_at: now,
386            updated_at: now,
387            archived_at: None,
388            deleted_at: None,
389        };
390        harness_store.add_harness(harness).await;
391
392        // Surface explicitly-added tools (via `.tool(...)`) as agent tool
393        // definitions so ReasonAtom returns them and ActAtom executes them
394        // (rather than treating them as unknown). Capability-provided tools are
395        // already surfaced through the capability registry.
396        let explicit_tool_definitions: Vec<crate::tool_types::ToolDefinition> =
397            self.tools.iter().map(|tool| tool.to_definition()).collect();
398
399        // Create agent
400        let agent_id = AgentId::new();
401        let agent = Agent {
402            public_id: agent_id,
403            internal_id: agent_id.uuid(),
404            name: "in-memory".to_string(),
405            display_name: Some(self.agent_name),
406            description: None,
407            system_prompt: self.system_prompt,
408            default_model_id: None,
409            harness_id,
410            default_version_id: None,
411            forked_from_agent_id: None,
412            forked_from_version_id: None,
413            root_agent_id: None,
414            tags: vec![],
415            capabilities: agent_capability_configs,
416            mcp_servers: Default::default(),
417            initial_files: vec![],
418            network_access: None,
419            max_iterations: None,
420            parallel_tool_calls: self.parallel_tool_calls,
421            tools: explicit_tool_definitions,
422            status: AgentStatus::Active,
423            created_at: now,
424            updated_at: now,
425            archived_at: None,
426            deleted_at: None,
427            usage: None,
428        };
429        agent_store.add_agent(agent).await;
430
431        // Create session
432        let session_id = SessionId::new();
433        let session = Session {
434            id: session_id,
435            workspace_id: crate::WorkspaceId::from_uuid((session_id).uuid()),
436            organization_id: crate::DEFAULT_ORG_PUBLIC_ID.to_string(),
437            harness_id,
438            agent_id: Some(agent_id),
439            agent_version_id: None,
440            agent_identity_id: None,
441            owner_principal_id: crate::PrincipalId::from_seed(1),
442            resolved_owner_user_id: None,
443            owner: None,
444            effective_owner: None,
445            title: Some("In-Memory Session".to_string()),
446            goal: None,
447            locale: None,
448            preview: None,
449            output_preview: None,
450            tags: vec![],
451            model_id: None,
452            capabilities: vec![],
453            tools: vec![],
454            mcp_servers: Default::default(),
455            system_prompt: None,
456            initial_files: vec![],
457            hints: None,
458            network_access: None,
459            max_iterations: None,
460            parallel_tool_calls: None,
461            status: SessionStatus::Started,
462            created_at: now,
463            updated_at: now,
464            started_at: None,
465            finished_at: None,
466            usage: None,
467            is_pinned: None,
468            active_schedule_count: None,
469            features: vec![],
470            parent_session_id: None,
471            forked_from_session_id: None,
472            forked_from_sequence: None,
473            blueprint_id: None,
474            blueprint_config: None,
475        };
476        session_store.add_session(session).await;
477
478        // Capture the configured model name before the if-let consumes self.model.
479        // Used below to resolve model-adaptive capabilities against the right variant.
480        let configured_model = self.model.as_ref().map(|m| m.model.clone());
481
482        // Create provider store and driver registry
483        let provider_store = InMemoryProviderStore::new();
484        let driver_registry =
485            if let (Some(model), Some(registry)) = (self.model, self.driver_registry) {
486                // Use provided model and driver registry
487                provider_store.set_default_model(model).await;
488                registry
489            } else {
490                // Use LlmSim (default or explicitly configured)
491                let config = self.llm_sim_config.unwrap_or_default();
492                let model = ResolvedModel {
493                    model: "llmsim-model".to_string(),
494                    provider_type: DriverId::LlmSim,
495                    api_key: Some("fake-key".to_string()),
496                    base_url: None,
497                    provider_metadata: None,
498                };
499                provider_store.set_default_model(model).await;
500
501                // Create the driver once and share it across calls.
502                // This ensures sequence-based responses work correctly
503                // because the Arc counters are shared.
504                let driver = LlmSimDriver::new(config);
505                let mut registry = DriverRegistry::new();
506                registry.register(DriverId::LlmSim, move |_config| Box::new(driver.clone()));
507                registry
508            };
509
510        // Build tool registry - include tools from capabilities. Resolve each
511        // capability against the configured model so model-adaptive capabilities
512        // (e.g. auto_tool_search) contribute the right variant for this harness.
513        let configured_model_ref = configured_model.as_deref();
514        let mut tool_builder = ToolRegistryBuilder::new();
515        for capability in &self.capabilities {
516            let effective: &dyn crate::Capability = capability
517                .resolve_for_model(configured_model_ref)
518                .unwrap_or_else(|| capability.as_ref());
519            for tool in effective.tools() {
520                tool_builder = tool_builder.tool_boxed(tool);
521            }
522        }
523
524        // Add explicit tools (can override capability tools)
525        for tool in self.tools {
526            tool_builder = tool_builder.tool_boxed(tool);
527        }
528        let tool_registry = tool_builder.build();
529
530        // Create capability registry with added capabilities
531        let mut capability_registry = CapabilityRegistry::new();
532        for capability in self.capabilities {
533            capability_registry.register_boxed(capability);
534        }
535
536        let input_atom = InputAtom::new(message_retriever.clone());
537        let mut reason_atom = ReasonAtom::new(
538            harness_store.clone(),
539            agent_store.clone(),
540            session_store.clone(),
541            message_retriever.clone(),
542            provider_store.clone(),
543            capability_registry,
544            driver_registry,
545            event_emitter.clone(),
546        );
547        let mut act_atom = ActAtom::new(tool_registry.clone(), event_emitter.clone())
548            .with_tool_registry(Arc::new(tool_registry.clone()));
549        if let Some(handle) = &self.reasoning_effort_handle {
550            reason_atom = reason_atom.with_reasoning_effort_handle(handle.clone());
551            act_atom = act_atom.with_reasoning_effort_handle(handle.clone());
552        }
553
554        Ok(InMemoryAgenticLoop {
555            harness_id,
556            agent_id,
557            session_id,
558            harness_store,
559            agent_store,
560            session_store,
561            message_retriever,
562            provider_store,
563            event_emitter,
564            tool_registry,
565            input_atom: Arc::new(input_atom),
566            reason_atom: Arc::new(reason_atom),
567            act_atom: Arc::new(act_atom),
568            max_iterations: self.max_iterations,
569            reasoning_effort_handle: self.reasoning_effort_handle,
570        })
571    }
572}
573
574// ============================================================================
575// InMemoryAgenticLoop
576// ============================================================================
577
578/// In-memory agentic loop for testing and prototyping
579///
580/// Bundles all in-memory stores and atoms into a convenient interface
581/// for running agent turns without external dependencies.
582///
583/// # Example
584///
585/// ```ignore
586/// use everruns_core::in_memory_loop::InMemoryAgenticLoop;
587///
588/// // Simple usage with simulated LLM
589/// let mut loop_runner = InMemoryAgenticLoop::builder()
590///     .system_prompt("You are a helpful assistant.")
591///     .with_simulated_response("Hello! I can help you with that.")
592///     .build()
593///     .await?;
594///
595/// let result = loop_runner.run_turn("Hi there!").await?;
596/// assert!(result.success);
597/// println!("Response: {}", result.response);
598///
599/// // With real LLM (requires API key)
600/// let mut loop_runner = InMemoryAgenticLoop::builder()
601///     .with_real_llm()
602///     .tool(MyCustomTool)
603///     .build()
604///     .await?;
605/// ```
606pub struct InMemoryAgenticLoop {
607    harness_id: HarnessId,
608    agent_id: AgentId,
609    session_id: SessionId,
610    #[allow(dead_code)]
611    harness_store: InMemoryHarnessStore,
612    #[allow(dead_code)]
613    agent_store: InMemoryAgentStore,
614    #[allow(dead_code)]
615    session_store: InMemorySessionStore,
616    message_retriever: InMemoryMessageRetriever,
617    #[allow(dead_code)]
618    provider_store: InMemoryProviderStore,
619    event_emitter: BridgingEventEmitter,
620    tool_registry: ToolRegistry,
621    input_atom: Arc<InputAtom<InMemoryMessageRetriever>>,
622    reason_atom: Arc<ReasonAtom>,
623    act_atom: Arc<ActAtom<ToolRegistry, BridgingEventEmitter>>,
624    max_iterations: usize,
625    reasoning_effort_handle: Option<crate::traits::ReasoningEffortHandle>,
626}
627
628impl InMemoryAgenticLoop {
629    /// Create a new builder
630    pub fn builder() -> InMemoryAgenticLoopBuilder {
631        InMemoryAgenticLoopBuilder::new()
632    }
633
634    /// Get the agent ID
635    pub fn agent_id(&self) -> AgentId {
636        self.agent_id
637    }
638
639    /// Get the session ID
640    pub fn session_id(&self) -> SessionId {
641        self.session_id
642    }
643
644    /// Run a turn with the given user input
645    ///
646    /// Accepts either a string or an `InputMessage` for full control over
647    /// message options like reasoning effort.
648    ///
649    /// This executes the full agentic loop using the TurnStateMachine:
650    /// 1. Add user message
651    /// 2. Record input (InputAtom)
652    /// 3. Reason loop (ReasonAtom → ActAtom → repeat until done)
653    ///
654    /// The TurnStateMachine ensures consistent orchestration logic,
655    /// proper error handling (checking success flag), and turn ID management.
656    ///
657    /// # Examples
658    ///
659    /// ```ignore
660    /// // Simple string input
661    /// let result = runner.run_turn("Hello").await?;
662    ///
663    /// // Full InputMessage with controls
664    /// let input = InputMessage {
665    ///     role: MessageRole::User,
666    ///     content: vec![ContentPart::text("What is 2+2?")],
667    ///     controls: Some(Controls {
668    ///         model_id: None,
669    ///         reasoning: Some(ReasoningConfig { effort: Some("medium".into()) }),
670    ///     }),
671    ///     metadata: None,
672    ///     tags: vec![],
673    /// };
674    /// let result = runner.run_turn(input).await?;
675    /// ```
676    pub async fn run_turn(&self, input: impl Into<InputMessage>) -> Result<TurnResult> {
677        // The live effort override is turn-scoped: tools may set it for later
678        // LLM steps in this turn, but stale values must not override the next
679        // turn's message controls.
680        if let Some(handle) = &self.reasoning_effort_handle {
681            handle.set(None);
682        }
683
684        // Add user message
685        let message = self
686            .message_retriever
687            .add(self.session_id, input.into())
688            .await?;
689
690        // Create turn context and state machine
691        let turn_context = TurnContext::new(self.session_id, message.id, self.agent_id, 0);
692        let mut state_machine = TurnStateMachine::new(turn_context, self.max_iterations);
693
694        // Track last reason result for ActAtom
695        let mut last_reason_result: Option<crate::atoms::ReasonResult> = None;
696        // Track response_id from last reason call for chaining
697        let mut previous_response_id: Option<String> = None;
698
699        // Execute the turn using the state machine
700        loop {
701            match state_machine.next_action() {
702                TurnAction::ExecuteInput => {
703                    let base_context = AtomContext::new(
704                        state_machine.context().session_id,
705                        state_machine.context().turn_id,
706                        state_machine.context().input_message_id,
707                    );
708                    self.input_atom
709                        .execute(InputAtomInput {
710                            context: base_context,
711                        })
712                        .await?;
713                    state_machine.on_input_completed();
714                }
715
716                TurnAction::ExecuteReason => {
717                    let base_context = AtomContext::new(
718                        state_machine.context().session_id,
719                        state_machine.context().turn_id,
720                        state_machine.context().input_message_id,
721                    );
722                    let reason_result = self
723                        .reason_atom
724                        .execute(ReasonInput {
725                            context: base_context.next_exec(),
726                            harness_id: self.harness_id,
727                            agent_id: Some(self.agent_id),
728                            org_id: 0,
729                            mcp_tool_definitions: vec![],
730                            previous_response_id: previous_response_id.take(),
731                            iteration: state_machine.current_iteration() as u32 + 1,
732                        })
733                        .await?;
734
735                    let tool_call_count = reason_result.tool_calls.len();
736                    previous_response_id = reason_result.response_id.clone();
737                    // In-memory loop has no signal mechanism, so
738                    // has_pending_user_messages is always false.
739                    state_machine.on_reason_completed(
740                        reason_result.text.clone(),
741                        reason_result.has_tool_calls,
742                        tool_call_count,
743                        reason_result.success,
744                        reason_result.error.clone(),
745                        false,
746                    );
747
748                    // Store for ActAtom if needed
749                    if reason_result.has_tool_calls {
750                        last_reason_result = Some(reason_result);
751                    }
752                }
753
754                TurnAction::ExecuteAct => {
755                    let reason_result = last_reason_result
756                        .take()
757                        .expect("ExecuteAct requires prior ReasonResult with tool calls");
758                    let base_context = AtomContext::new(
759                        state_machine.context().session_id,
760                        state_machine.context().turn_id,
761                        state_machine.context().input_message_id,
762                    );
763                    self.act_atom
764                        .execute(ActInput {
765                            org_id: Some(0),
766                            context: base_context.next_exec(),
767                            harness_id: self.harness_id,
768                            agent_id: Some(self.agent_id),
769                            tool_calls: reason_result.tool_calls,
770                            tool_definitions: reason_result.tool_definitions,
771                            locale: reason_result.locale,
772                            blueprint_id: None,
773                            network_access: reason_result.network_access,
774                            // Request-level parallel tool calling preference,
775                            // carried from agent config through reason (EVE-598).
776                            parallel_tool_calls: reason_result.parallel_tool_calls,
777                        })
778                        .await?;
779                    state_machine.on_act_completed();
780                }
781
782                TurnAction::Complete(outcome) => {
783                    return Ok(TurnResult::from_outcome(
784                        outcome,
785                        state_machine.context().turn_id,
786                    ));
787                }
788            }
789        }
790    }
791
792    /// Run multiple turns in sequence
793    pub async fn run_conversation(&self, messages: &[&str]) -> Result<Vec<TurnResult>> {
794        let mut results = Vec::with_capacity(messages.len());
795        for msg in messages {
796            results.push(self.run_turn(*msg).await?);
797        }
798        Ok(results)
799    }
800
801    /// Get all messages in the session
802    pub async fn messages(&self) -> Result<Vec<Message>> {
803        self.message_retriever.load(self.session_id).await
804    }
805
806    /// Get all emitted events
807    pub async fn events(&self) -> Vec<Event> {
808        self.event_emitter.events().await
809    }
810
811    /// Get events of a specific type
812    pub async fn events_by_type(&self, event_type: &str) -> Vec<Event> {
813        self.event_emitter.events_by_type(event_type).await
814    }
815
816    /// Get the count of messages
817    pub async fn message_count(&self) -> Result<usize> {
818        self.message_retriever.count(self.session_id).await
819    }
820
821    /// Get the count of events
822    pub async fn event_count(&self) -> usize {
823        self.event_emitter.event_count().await
824    }
825
826    /// Clear all events (useful between tests)
827    pub async fn clear_events(&self) {
828        self.event_emitter.clear().await;
829    }
830
831    /// Clear all messages (starts a fresh conversation)
832    pub async fn clear_messages(&self) {
833        self.message_retriever.clear_session(self.session_id).await;
834    }
835
836    /// Reset the loop (clear messages and events)
837    pub async fn reset(&self) {
838        self.clear_messages().await;
839        self.clear_events().await;
840    }
841
842    /// Get conversation as a formatted string
843    pub async fn conversation_string(&self) -> Result<String> {
844        let messages = self.messages().await?;
845        let mut result = String::new();
846        for msg in messages {
847            let role = format!("{:?}", msg.role);
848            let text = msg.text().unwrap_or("[non-text content]");
849            result.push_str(&format!("[{}] {}\n", role, text));
850        }
851        Ok(result)
852    }
853
854    /// Access the message retriever directly
855    pub fn message_retriever(&self) -> &InMemoryMessageRetriever {
856        &self.message_retriever
857    }
858
859    /// Access the tool registry directly
860    pub fn tool_registry(&self) -> &ToolRegistry {
861        &self.tool_registry
862    }
863}
864
865// ============================================================================
866// Quick constructors
867// ============================================================================
868
869impl InMemoryAgenticLoop {
870    /// Create a simple loop with a fixed simulated response
871    ///
872    /// # Example
873    ///
874    /// ```ignore
875    /// let runner = InMemoryAgenticLoop::with_fixed_response("Hello!").await?;
876    /// let result = runner.run_turn("Hi").await?;
877    /// assert_eq!(result.response, "Hello!");
878    /// ```
879    pub async fn with_fixed_response(response: impl Into<String>) -> Result<Self> {
880        Self::builder()
881            .with_simulated_response(response)
882            .build()
883            .await
884    }
885
886    /// Create a loop that echoes user input
887    ///
888    /// # Example
889    ///
890    /// ```ignore
891    /// let runner = InMemoryAgenticLoop::with_echo().await?;
892    /// let result = runner.run_turn("Hello").await?;
893    /// assert!(result.response.contains("Hello"));
894    /// ```
895    pub async fn with_echo() -> Result<Self> {
896        Self::builder()
897            .with_llm_sim(LlmSimConfig::echo())
898            .build()
899            .await
900    }
901
902    /// Create a loop with sequence of responses
903    ///
904    /// # Example
905    ///
906    /// ```ignore
907    /// let runner = InMemoryAgenticLoop::with_sequence(vec![
908    ///     "First response",
909    ///     "Second response",
910    /// ]).await?;
911    ///
912    /// let r1 = runner.run_turn("msg1").await?;
913    /// let r2 = runner.run_turn("msg2").await?;
914    /// assert_eq!(r1.response, "First response");
915    /// assert_eq!(r2.response, "Second response");
916    /// ```
917    pub async fn with_sequence(responses: Vec<impl Into<String>>) -> Result<Self> {
918        let responses: Vec<String> = responses.into_iter().map(|s| s.into()).collect();
919        Self::builder()
920            .with_llm_sim(LlmSimConfig::sequence(responses))
921            .build()
922            .await
923    }
924
925    /// Create a loop with tool call simulation
926    ///
927    /// # Example
928    ///
929    /// ```ignore
930    /// use everruns_core::ToolCall;
931    /// use serde_json::json;
932    ///
933    /// let tool_call = ToolCall {
934    ///     id: "call_1".to_string(),
935    ///     name: "get_weather".to_string(),
936    ///     arguments: json!({"city": "NYC"}),
937    /// };
938    ///
939    /// let runner = InMemoryAgenticLoop::with_tool_calls(
940    ///     "Let me check that.",
941    ///     vec![tool_call],
942    /// ).await?;
943    /// ```
944    pub async fn with_tool_calls(
945        response: impl Into<String>,
946        tool_calls: Vec<ToolCall>,
947    ) -> Result<Self> {
948        Self::builder()
949            .with_llm_sim(LlmSimConfig::fixed(response).with_tool_calls(tool_calls))
950            .build()
951            .await
952    }
953}
954
955// ============================================================================
956// Tests
957// ============================================================================
958
959#[cfg(test)]
960mod tests {
961    use super::*;
962
963    #[tokio::test]
964    async fn test_simple_turn() {
965        let runner = InMemoryAgenticLoop::with_fixed_response("Hello from the assistant!")
966            .await
967            .unwrap();
968
969        let result = runner.run_turn("Hi there").await.unwrap();
970
971        assert!(result.success);
972        assert_eq!(result.response, "Hello from the assistant!");
973        assert_eq!(result.iterations, 1);
974        assert_eq!(result.tool_calls_count, 0);
975    }
976
977    #[tokio::test]
978    async fn test_echo_turn() {
979        let runner = InMemoryAgenticLoop::with_echo().await.unwrap();
980
981        let result = runner.run_turn("Test message").await.unwrap();
982
983        assert!(result.success);
984        assert!(result.response.contains("Test message"));
985    }
986
987    #[tokio::test]
988    async fn test_sequence_turns() {
989        let runner = InMemoryAgenticLoop::with_sequence(vec!["First", "Second", "Third"])
990            .await
991            .unwrap();
992
993        let r1 = runner.run_turn("msg1").await.unwrap();
994        let r2 = runner.run_turn("msg2").await.unwrap();
995        let r3 = runner.run_turn("msg3").await.unwrap();
996
997        assert_eq!(r1.response, "First");
998        assert_eq!(r2.response, "Second");
999        assert_eq!(r3.response, "Third");
1000    }
1001
1002    #[tokio::test]
1003    async fn test_conversation() {
1004        let runner = InMemoryAgenticLoop::with_sequence(vec!["Hello!", "How can I help?"])
1005            .await
1006            .unwrap();
1007
1008        let results = runner
1009            .run_conversation(&["Hi", "I need help"])
1010            .await
1011            .unwrap();
1012
1013        assert_eq!(results.len(), 2);
1014        assert_eq!(results[0].response, "Hello!");
1015        assert_eq!(results[1].response, "How can I help?");
1016
1017        // Check message history
1018        let messages = runner.messages().await.unwrap();
1019        assert_eq!(messages.len(), 4); // 2 user + 2 assistant
1020    }
1021
1022    #[tokio::test]
1023    async fn test_events_captured() {
1024        let runner = InMemoryAgenticLoop::with_fixed_response("Response")
1025            .await
1026            .unwrap();
1027
1028        runner.run_turn("Test").await.unwrap();
1029
1030        let events = runner.events().await;
1031        assert!(!events.is_empty());
1032
1033        // Should have reason.* events (input.message is emitted by API layer, not InputAtom)
1034        let reason_events = runner.events_by_type("reason.started").await;
1035        assert_eq!(reason_events.len(), 1);
1036    }
1037
1038    #[tokio::test]
1039    async fn test_reset() {
1040        let runner = InMemoryAgenticLoop::with_fixed_response("Response")
1041            .await
1042            .unwrap();
1043
1044        runner.run_turn("Test").await.unwrap();
1045        assert!(runner.message_count().await.unwrap() > 0);
1046        assert!(runner.event_count().await > 0);
1047
1048        runner.reset().await;
1049        assert_eq!(runner.message_count().await.unwrap(), 0);
1050        assert_eq!(runner.event_count().await, 0);
1051    }
1052
1053    #[tokio::test]
1054    async fn test_builder_with_custom_config() {
1055        let runner = InMemoryAgenticLoop::builder()
1056            .agent_name("Custom Agent")
1057            .system_prompt("You are a custom assistant.")
1058            .with_simulated_response("Custom response")
1059            .max_iterations(5)
1060            .build()
1061            .await
1062            .unwrap();
1063
1064        let result = runner.run_turn("Test").await.unwrap();
1065        assert_eq!(result.response, "Custom response");
1066    }
1067
1068    #[tokio::test]
1069    async fn test_conversation_string() {
1070        let runner = InMemoryAgenticLoop::with_fixed_response("Hello!")
1071            .await
1072            .unwrap();
1073
1074        runner.run_turn("Hi").await.unwrap();
1075
1076        let conv = runner.conversation_string().await.unwrap();
1077        assert!(conv.contains("[User]"));
1078        assert!(conv.contains("[Agent]"));
1079        assert!(conv.contains("Hi"));
1080        assert!(conv.contains("Hello!"));
1081    }
1082}