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::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}