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