Skip to main content

agentd/runtime/
reload.rs

1// SPDX-License-Identifier: AGPL-3.0-only
2//! **Hot reload** of the configuration: SIGHUP or `lifecycle.watch_config`
3//! re-merges the files and re-validates. The reload is all-or-nothing — if any
4//! restart-only path changed the whole reload is refused as
5//! `restart_required` and the running configuration stays, so the daemon never
6//! ends up half on one configuration and half on another.
7//!
8//! The reloadable partition applies at the loop's quiesce boundary. The flat
9//! child tree makes most of it trivial: every turn worker is spawned fresh
10//! from the live settings, so a new intelligence endpoint, model, instruction,
11//! budget, tool override or workflow definition takes effect for the next unit
12//! of work without touching the units already in flight. Live workflow runs
13//! keep the definition they started with, pinned by hash, so a run cannot
14//! change shape halfway through.
15
16use 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    /// SIGHUP / watcher: reload, diff, apply (or refuse).
27    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        // Audit the reload: reconfiguring a running daemon is an operator
37        // action, and a refused reload is as worth recording as an applied one.
38        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        // Intelligence is hot-swappable: workers in flight keep dialing the
106        // endpoint they were spawned with, the next spawned worker uses this.
107        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        // Budgets: new windows, counters carried over.
123        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        // Instruction (static text; a resource instruction re-subscribes).
131        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        // MCP servers: connect added, drop removed (re-handshake).
177        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        // Registry (overrides/disabled/tools) — always rebuilt when tools/mcp/knowledge/search changed.
236        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                    // The rebuild failed, so the daemon stays on the tool
266                    // configuration it is already running: put back the tools
267                    // settings to match the registry that is still installed,
268                    // or the two would disagree about what is callable.
269                    self.settings.tools = old.tools.clone();
270                    return Err(ReloadRefused::Invalid(errs));
271                }
272            }
273        }
274        // Skills sources — the config section, or the instruction's inline
275        // `:::skill` definitions (they live on `agent`, but they land in this
276        // catalogue).
277        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        // Workflows: reload definitions. Retirement (runtime::retire) gives
303        // every old version the same exit — unsubscribe what nothing else
304        // wants, pin for live runs, apply its own `unload:` policy — whether
305        // it was removed outright or replaced by a new hash.
306        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; // the running set stays authoritative
310                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
346/// Why a reload did not apply.
347enum ReloadRefused {
348    Invalid(Vec<String>),
349    RestartRequired(Vec<String>),
350}