Skip to main content

vv_agent/runtime/engine/
controls.rs

1use std::collections::{BTreeMap, VecDeque};
2use std::path::PathBuf;
3use std::sync::{Arc, Mutex};
4
5use serde_json::Value;
6
7use crate::budget::{BudgetUsageSnapshot, HostCostMeter, RunBudgetLimits};
8use crate::events::RunEvent;
9use crate::model::ModelProvider;
10use crate::runtime::cancellation::CancellationToken;
11use crate::runtime::checkpoint_resume::CheckpointController;
12use crate::runtime::context::ExecutionContext;
13use crate::runtime::sub_task_manager::SubTaskManager;
14use crate::types::{CycleRecord, Message, ModelCallRecord};
15use crate::workspace::WorkspaceBackend;
16use crate::{RunConfig, RunContext};
17
18pub type RunEventHandler = Arc<dyn Fn(&RunEvent) + Send + Sync + 'static>;
19pub type BeforeCycleMessageProvider =
20    Arc<dyn Fn(u32, &[Message], &BTreeMap<String, Value>) -> Vec<Message> + Send + Sync + 'static>;
21pub type InterruptionMessageProvider = Arc<dyn Fn() -> Vec<Message> + Send + Sync + 'static>;
22
23#[doc(hidden)]
24#[derive(Clone)]
25pub struct CheckpointRuntimeControl {
26    controller: CheckpointController,
27}
28
29impl CheckpointRuntimeControl {
30    pub(crate) fn new(controller: CheckpointController) -> Self {
31        Self { controller }
32    }
33
34    pub(crate) fn controller(&self) -> &CheckpointController {
35        &self.controller
36    }
37
38    pub(crate) fn into_controller(self) -> CheckpointController {
39        self.controller
40    }
41}
42
43#[derive(Clone, Default)]
44pub struct RuntimeRunControls {
45    pub event_handler: Option<RunEventHandler>,
46    pub before_cycle_messages: Option<BeforeCycleMessageProvider>,
47    pub interruption_messages: Option<InterruptionMessageProvider>,
48    pub steering_queue: Option<Arc<Mutex<VecDeque<String>>>>,
49    pub cancellation_token: Option<CancellationToken>,
50    pub execution_context: Option<ExecutionContext>,
51    pub workspace: Option<PathBuf>,
52    pub workspace_backend: Option<Arc<dyn WorkspaceBackend>>,
53    pub model_provider: Option<Arc<dyn ModelProvider>>,
54    pub run_context: Option<RunContext>,
55    pub sub_task_manager: Option<SubTaskManager>,
56    pub budget_limits: Option<RunBudgetLimits>,
57    pub host_cost_meter: Option<Arc<dyn HostCostMeter>>,
58    #[doc(hidden)]
59    pub background_parent_run_config: Option<RunConfig>,
60    #[doc(hidden)]
61    pub initial_messages: Option<Vec<Message>>,
62    #[doc(hidden)]
63    pub initial_shared_state: Option<BTreeMap<String, Value>>,
64    #[doc(hidden)]
65    pub initial_cycles: Option<Vec<CycleRecord>>,
66    #[doc(hidden)]
67    pub initial_model_calls: Option<Vec<ModelCallRecord>>,
68    #[doc(hidden)]
69    pub cycle_index_start: Option<u32>,
70    #[doc(hidden)]
71    pub cycle_count: Option<u32>,
72    #[doc(hidden)]
73    pub initial_budget_usage: Option<BudgetUsageSnapshot>,
74    #[doc(hidden)]
75    pub defer_terminal_on_max_cycles: bool,
76    #[doc(hidden)]
77    pub checkpoint_controller: Option<CheckpointRuntimeControl>,
78}
79
80impl RuntimeRunControls {
81    pub(in crate::runtime::engine) fn effective_cancellation_token(
82        &self,
83    ) -> Option<CancellationToken> {
84        self.cancellation_token.clone().or_else(|| {
85            self.execution_context
86                .as_ref()
87                .and_then(|context| context.cancellation_token.clone())
88        })
89    }
90
91    pub(in crate::runtime::engine) fn effective_event_handler(&self) -> Option<RunEventHandler> {
92        self.execution_context
93            .as_ref()
94            .and_then(|context| context.event_handler.clone())
95            .or_else(|| self.event_handler.clone())
96    }
97
98    pub(crate) fn effective_checkpoint_controller(&self) -> Option<&CheckpointController> {
99        self.checkpoint_controller
100            .as_ref()
101            .map(CheckpointRuntimeControl::controller)
102    }
103}