1use super::reactor::Runtime;
17use crate::config::v2 as cfg;
18use crate::governor::Governor;
19use crate::registry::{Registry, ServerTools};
20use crate::state::now_ms;
21use serde_json::json;
22use std::sync::Arc;
23use std::time::Duration;
24
25impl Runtime {
26 pub(crate) fn on_reload_requested(&mut self) {
28 let trigger = if crate::signals::take_reload_was_watch() {
29 "watch"
30 } else {
31 "sighup"
32 };
33 crate::signals::set_reloading(true);
34 let outcome = self.reload_inner();
35 crate::signals::set_reloading(false);
36 let (label, atarget) = match &outcome {
39 Ok(changed) => ("applied", json!({"trigger": trigger, "changed": changed})),
40 Err(ReloadRefused::Invalid(errs)) => {
41 ("invalid", json!({"trigger": trigger, "errors": errs}))
42 }
43 Err(ReloadRefused::RestartRequired(paths)) => (
44 "restart_required",
45 json!({"trigger": trigger, "paths": paths}),
46 ),
47 };
48 self.audit(crate::runtime::audit::AuditEvent {
49 action: "config.reload",
50 target: atarget,
51 outcome: label,
52 principal: Some("operator"),
53 role: Some("operator"),
54 request_id: None,
55 });
56 match outcome {
57 Ok(changed) => {
58 self.log.info(
59 "config.reloaded",
60 json!({"trigger": trigger, "changed": changed}),
61 );
62 crate::obs::metrics::record_config_reload("applied");
63 let mut generation = 0;
64 self.durable.manifest_update(|m| {
65 generation = m.lifecycle["config_generation"].as_u64().unwrap_or(0) + 1;
66 m.lifecycle["config_generation"] = json!(generation);
67 m.lifecycle["config_reloaded_at"] = json!(now_ms());
68 });
69 crate::obs::metrics::set_config_generation(generation);
70 }
71 Err(ReloadRefused::Invalid(errs)) => {
72 for e in &errs {
73 self.log.warn(
74 "config.reload.invalid",
75 json!({"trigger": trigger, "error": e}),
76 );
77 }
78 crate::obs::metrics::record_config_reload("invalid");
79 }
80 Err(ReloadRefused::RestartRequired(paths)) => {
81 self.log.warn(
82 "config.reload.restart_required",
83 json!({"trigger": trigger, "paths": paths}),
84 );
85 crate::obs::metrics::record_config_reload("restart_required");
86 }
87 }
88 }
89
90 fn reload_inner(&mut self) -> Result<Vec<&'static str>, ReloadRefused> {
91 let (loaded, _ask) = cfg::load(&self.args, &self.env)
92 .map_err(|e| ReloadRefused::Invalid(vec![format!("{e:?}")]))?;
93 let restart = cfg::restart_only_diff(&self.settings_doc, &loaded.doc);
94 if !restart.is_empty() {
95 return Err(ReloadRefused::RestartRequired(restart));
96 }
97 for w in &loaded.warnings {
98 self.log.warn("config.warning", json!({"warning": w}));
99 }
100 let new = loaded.settings;
101 let old = std::mem::replace(&mut self.settings, new.clone());
102 self.settings_doc = loaded.doc;
103 let mut changed = Vec::new();
104
105 if old.intelligence.endpoints != new.intelligence.endpoints
108 || old.intelligence.model != new.intelligence.model
109 || old.intelligence.token != new.intelligence.token
110 || old.intelligence.token_file != new.intelligence.token_file
111 {
112 self.intel_uri = new.intelligence.endpoint_list().unwrap_or_default();
113 self.model = new.intelligence.model.clone().unwrap_or_default();
114 let env = self.env.clone();
115 let envmap = move |k: &str| env.iter().find(|(n, _)| n == k).map(|(_, v)| v.clone());
116 match super::resolve_intel_token(&new, &envmap) {
117 Ok(t) => self.intel_token = t,
118 Err(e) => self.log.warn("config.reload.token", json!({"err": e})),
119 }
120 changed.push("intelligence");
121 }
122 if old.intelligence.budget != new.intelligence.budget {
124 let counters = self.governor.to_value();
125 let mut g = Governor::new(&new.intelligence.budget);
126 g.restore(&counters, now_ms());
127 self.governor = g;
128 changed.push("intelligence.budget");
129 }
130 if old.agent.instruction != new.agent.instruction {
132 match new.agent.instruction.clone() {
133 Some(t) if cfg::looks_like_resource_uri(&t) => {
134 if let Err(e) = self.subscribe_instruction(&t) {
135 self.log
136 .warn("instruction.subscribe.fail", json!({"uri": t, "err": e}));
137 }
138 }
139 Some(t) => {
140 self.instruction = super::reactor::Instruction {
141 text: t,
142 source: "static",
143 uri: None,
144 server: None,
145 version: self.instruction.version + 1,
146 };
147 }
148 None => {
149 self.instruction = super::reactor::Instruction {
150 text: String::new(),
151 source: "static",
152 uri: None,
153 server: None,
154 version: self.instruction.version + 1,
155 };
156 }
157 }
158 if new
159 .agent
160 .wake_on()
161 .contains(&cfg::WakeEvent::InstructionUpdated)
162 {
163 self.note_root("instruction.updated: the configuration changed the instruction; re-read it with instruction.read".into());
164 }
165 changed.push("agent.instruction");
166 }
167 if old.agent.preflight != new.agent.preflight
168 || old.agent.wake_on != new.agent.wake_on
169 || old.agent.tools != new.agent.tools
170 || old.agent.max_parallel_turns != new.agent.max_parallel_turns
171 || old.agent.on_workflow_finished != new.agent.on_workflow_finished
172 || old.agent.conversation_budget != new.agent.conversation_budget
173 {
174 changed.push("agent");
175 }
176 if old.mcp != new.mcp {
178 let keep: Vec<String> = new.mcp.servers.iter().map(|s| s.name.clone()).collect();
179 let removed: Vec<String> = self
180 .mcp
181 .keys()
182 .filter(|k| !keep.contains(k))
183 .cloned()
184 .collect();
185 for r in &removed {
186 self.mcp.remove(r);
187 self.mcp_specs.remove(r);
188 self.skills.forget_server(r);
189 self.log
190 .info("mcp.disconnect", json!({"server": r, "reason": "reload"}));
191 }
192 let timeout = new
193 .mcp
194 .default_timeout
195 .map(|d| d.0)
196 .unwrap_or(Duration::from_secs(60));
197 for s in &new.mcp.servers {
198 let spec = match s.to_spec() {
199 Ok(sp) => sp,
200 Err(e) => {
201 self.log
202 .warn("mcp.spec.invalid", json!({"server": s.name, "err": e}));
203 continue;
204 }
205 };
206 let same = self.mcp_specs.get(&s.name).is_some_and(|old| {
207 old.endpoint == spec.endpoint
208 && old.headers == spec.headers
209 && old.aauth == spec.aauth
210 });
211 if same && self.mcp.contains_key(&s.name) {
212 self.mcp_specs.insert(s.name.clone(), spec);
213 continue;
214 }
215 match crate::mcp::from_spec(&spec, s.timeout.map(|d| d.0).unwrap_or(timeout))
216 .and_then(|mut c| c.initialize().map(|()| c))
217 {
218 Ok(mut c) => {
219 c.set_tool_meta(
220 json!({"agent/run_id": self.run_id, "agent/instance": self.instance}),
221 );
222 self.log
223 .info("mcp.connect", json!({"server": s.name, "reason": "reload"}));
224 self.mcp.insert(s.name.clone(), Arc::new(c));
225 }
226 Err(e) => self.log.warn(
227 "mcp.connect.fail",
228 json!({"server": s.name, "err": e.to_string()}),
229 ),
230 }
231 self.mcp_specs.insert(s.name.clone(), spec);
232 }
233 changed.push("mcp");
234 }
235 if old.tools != new.tools
237 || old.mcp != new.mcp
238 || old.knowledge != new.knowledge
239 || old.search != new.search
240 {
241 let server_tools: Vec<ServerTools> = new
242 .mcp
243 .servers
244 .iter()
245 .filter_map(|s| {
246 let c = self.mcp.get(&s.name)?;
247 Some(ServerTools {
248 name: s.name.clone(),
249 ns: s.ns.clone(),
250 tags: self
251 .mcp_specs
252 .get(&s.name)
253 .map(|sp| sp.tags.clone())
254 .unwrap_or_default(),
255 tools: c.list_tools().unwrap_or_default(),
256 })
257 })
258 .collect();
259 match Registry::build(&new, &server_tools) {
260 Ok(r) => {
261 self.registry = r;
262 changed.push("tools");
263 }
264 Err(errs) => {
265 self.settings.tools = old.tools.clone();
270 return Err(ReloadRefused::Invalid(errs));
271 }
272 }
273 }
274 if old.skills != new.skills || old.agent.inline_skills != new.agent.inline_skills {
278 let mut cat = crate::context::skills::Catalogue::new(
279 new.skills
280 .reference_prefix
281 .as_deref()
282 .unwrap_or(crate::context::skills::DEFAULT_PREFIX),
283 new.skills.max_bytes.unwrap_or(32_768) as usize,
284 );
285 for src in &new.skills.sources {
286 if let Some(c) = self.mcp.get(&src.server) {
287 let mode = match src.discover {
288 cfg::Discover::Prompts => crate::context::skills::Discover::Prompts,
289 cfg::Discover::Resources => crate::context::skills::Discover::Resources,
290 cfg::Discover::Auto => crate::context::skills::Discover::Auto,
291 };
292 cat.discover(&**c, mode, src.filter.as_deref());
293 }
294 }
295 if let Some(dir) = &new.skills.dir {
296 cat.add_dir(std::path::Path::new(dir));
297 }
298 cat.add_inline(&new.agent.inline_skills);
299 self.skills = cat;
300 changed.push("skills");
301 }
302 if old.workflows != new.workflows {
307 let previous = std::mem::take(&mut self.workflows);
308 if let Err(errs) = self.load_workflows() {
309 self.workflows = previous; return Err(ReloadRefused::Invalid(errs));
311 }
312 for (name, wf) in &previous {
313 let survives = self
314 .workflows
315 .get(name)
316 .is_some_and(|new_wf| new_wf.hash == wf.hash);
317 if survives {
318 continue;
319 }
320 let reason = if self.workflows.contains_key(name) {
321 "replaced"
322 } else {
323 "removed"
324 };
325 self.retire_workflow(wf, reason);
326 }
327 self.arm_workflows();
328 changed.push("workflows");
329 }
330 if old.limits != new.limits
331 || old.lifecycle.idle_grace != new.lifecycle.idle_grace
332 || old.observability.log_level != new.observability.log_level
333 || old.observability.log_content != new.observability.log_content
334 || old.memory != new.memory
335 || old.context != new.context
336 {
337 changed.push("limits/lifecycle/observability/memory/context");
338 }
339 if changed.is_empty() {
340 changed.push("nothing");
341 }
342 Ok(changed)
343 }
344}
345
346enum ReloadRefused {
348 Invalid(Vec<String>),
349 RestartRequired(Vec<String>),
350}