vv_agent/runtime/engine/
controls.rs1use 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}