1use crate::nodes::Target;
20use sha2::{Digest, Sha256};
21use std::io::Read;
22use std::path::{Path, PathBuf};
23use std::process::{Command, Stdio};
24use std::time::{Duration, Instant};
25
26pub const KEEP_VERSIONS: usize = 3;
27pub const COMMAND_TIMEOUT: Duration = Duration::from_secs(30);
29
30pub fn sha256_hex(bytes: &[u8]) -> String {
31 hex::encode(Sha256::digest(bytes))
32}
33
34pub fn hash_file(path: &Path) -> Result<Option<String>, String> {
36 match std::fs::read(path) {
37 Ok(b) => Ok(Some(sha256_hex(&b))),
38 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(None),
39 Err(e) => Err(format!("{}: {e}", path.display())),
40 }
41}
42
43fn run(cmd: &str, content: &str, timeout: Duration) -> Result<(), String> {
46 #[cfg(windows)]
47 let mut c = {
48 let mut c = Command::new("cmd");
49 c.args(["/C", cmd]);
50 c
51 };
52 #[cfg(not(windows))]
53 let mut c = {
54 let mut c = Command::new("sh");
55 c.args(["-c", cmd]);
56 c
57 };
58 let mut child = c
59 .stdin(Stdio::null())
60 .stdout(Stdio::piped())
61 .stderr(Stdio::piped())
62 .spawn()
63 .map_err(|e| format!("could not start '{cmd}': {e}"))?;
64 let start = Instant::now();
65 let status = loop {
66 match child.try_wait().map_err(|e| e.to_string())? {
67 Some(s) => break s,
68 None if start.elapsed() > timeout => {
69 let _ = child.kill();
70 let _ = child.wait();
71 return Err(format!(
72 "'{cmd}' did not finish in {}s",
73 timeout.as_secs().max(1)
74 ));
75 }
76 None => std::thread::sleep(Duration::from_millis(25)),
77 }
78 };
79 if status.success() {
80 return Ok(());
81 }
82 let mut text = String::new();
83 if let Some(mut e) = child.stderr.take() {
84 let _ = e.by_ref().take(8192).read_to_string(&mut text);
85 }
86 if text.trim().is_empty() {
87 if let Some(mut o) = child.stdout.take() {
88 let _ = o.by_ref().take(8192).read_to_string(&mut text);
89 }
90 }
91 Err(format!(
92 "'{cmd}' exited {}: {}",
93 status.code().map_or("by signal".into(), |c| c.to_string()),
94 scrub(&text, content)
95 ))
96}
97
98pub fn scrub(output: &str, content: &str) -> String {
103 let file_lines: std::collections::HashSet<&str> = content
104 .lines()
105 .map(str::trim)
106 .filter(|l| l.len() >= 8)
107 .collect();
108 let lines: Vec<&str> = output
109 .lines()
110 .map(str::trim)
111 .filter(|l| !l.is_empty())
112 .collect();
113 let tail = &lines[lines.len().saturating_sub(3)..];
114 tail.iter()
115 .map(|l| {
116 if file_lines.contains(l) || file_lines.iter().any(|f| l.contains(f)) {
117 "[a line of the file elided]".to_string()
118 } else {
119 l.chars().take(200).collect()
120 }
121 })
122 .collect::<Vec<_>>()
123 .join(" | ")
124}
125
126fn versions_dir(state_dir: &Path, id: &str) -> PathBuf {
127 state_dir.join("versions").join(id)
128}
129
130fn write_atomic(
131 path: &Path,
132 bytes: &[u8],
133 mode_from: Option<&std::fs::Metadata>,
134) -> Result<(), String> {
135 use std::io::Write;
136 let dir = path.parent().ok_or("target has no parent directory")?;
137 let name = path.file_name().ok_or("target has no file name")?;
138 let tmp = dir.join(format!(".{}.unv-node.tmp", name.to_string_lossy()));
139 let mut opts = std::fs::OpenOptions::new();
140 opts.write(true).create(true).truncate(true);
141 #[cfg(unix)]
142 {
143 use std::os::unix::fs::{MetadataExt, OpenOptionsExt};
144 opts.mode(mode_from.map_or(0o600, |m| m.mode() & 0o7777));
146 }
147 #[cfg(not(unix))]
148 let _ = mode_from;
149 let mut f = opts
150 .open(&tmp)
151 .map_err(|e| format!("cannot write beside {}: {e}", path.display()))?;
152 f.write_all(bytes).map_err(|e| e.to_string())?;
153 f.sync_all().map_err(|e| e.to_string())?;
154 drop(f);
155 std::fs::rename(&tmp, path).map_err(|e| {
156 let _ = std::fs::remove_file(&tmp);
157 format!("cannot replace {}: {e}", path.display())
158 })?;
159 #[cfg(unix)]
160 if let Ok(d) = std::fs::File::open(dir) {
161 let _ = d.sync_all();
162 }
163 Ok(())
164}
165
166fn restore(
167 path: &Path,
168 previous: &Option<Vec<u8>>,
169 meta: Option<&std::fs::Metadata>,
170) -> Result<(), String> {
171 match previous {
172 Some(b) => write_atomic(path, b, meta),
173 None => std::fs::remove_file(path).map_err(|e| e.to_string()),
174 }
175}
176
177fn keep_version(state_dir: &Path, id: &str, bytes: &[u8]) -> Result<(), String> {
178 let dir = versions_dir(state_dir, id);
179 std::fs::create_dir_all(&dir).map_err(|e| e.to_string())?;
180 let secs = std::time::SystemTime::now()
181 .duration_since(std::time::UNIX_EPOCH)
182 .map_or(0, |d| d.as_secs());
183 let name = format!("{secs:012}-{}", &sha256_hex(bytes)[..8]);
184 let p = dir.join(name);
185 write_atomic(&p, bytes, None)?;
188 let mut all: Vec<_> = std::fs::read_dir(&dir)
189 .map_err(|e| e.to_string())?
190 .flatten()
191 .map(|e| e.path())
192 .filter(|p| {
193 !p.file_name()
194 .is_some_and(|n| n.to_string_lossy().starts_with('.'))
195 })
196 .collect();
197 all.sort();
198 while all.len() > KEEP_VERSIONS {
199 let _ = std::fs::remove_file(all.remove(0));
200 }
201 Ok(())
202}
203
204pub fn kept_versions(state_dir: &Path, id: &str) -> Vec<PathBuf> {
206 let mut v: Vec<PathBuf> = std::fs::read_dir(versions_dir(state_dir, id))
207 .map(|d| d.flatten().map(|e| e.path()).collect())
208 .unwrap_or_default();
209 v.sort();
210 v
211}
212
213#[derive(Debug, PartialEq)]
214pub enum Applied {
215 Unchanged,
217 Written,
218}
219
220pub fn apply(
222 target: &Target,
223 content: &[u8],
224 expected_sha: &str,
225 state_dir: &Path,
226) -> Result<Applied, String> {
227 if target.mode != "push" || !target.apply {
228 return Err(format!(
229 "target '{}' is not set to apply; the node's config decides that, not the hub",
230 target.id
231 ));
232 }
233 if sha256_hex(content) != expected_sha {
234 return Err(
235 "content does not match the hash the hub announced; nothing was written".into(),
236 );
237 }
238 let text = String::from_utf8_lossy(content).into_owned();
239 let path = &target.path;
240 let meta = std::fs::metadata(path).ok();
241 let previous = match std::fs::read(path) {
242 Ok(b) => Some(b),
243 Err(e) if e.kind() == std::io::ErrorKind::NotFound => None,
244 Err(e) => return Err(format!("{}: {e}", path.display())),
245 };
246 if previous.as_deref() == Some(content) {
247 return Ok(Applied::Unchanged);
248 }
249 if let Some(p) = &previous {
250 keep_version(state_dir, &target.id, p)?;
251 }
252 write_atomic(path, content, meta.as_ref())?;
253
254 if let Some(v) = &target.validate {
255 if let Err(e) = run(v, &text, COMMAND_TIMEOUT) {
256 return Err(match restore(path, &previous, meta.as_ref()) {
257 Ok(()) => format!("validate failed, previous file restored: {e}"),
258 Err(r) => {
259 format!("validate failed AND restoring the previous file failed ({r}): {e}")
260 }
261 });
262 }
263 }
264 if let Some(r) = &target.reload {
265 if let Err(e) = run(r, &text, COMMAND_TIMEOUT) {
266 return Err(match restore(path, &previous, meta.as_ref()) {
267 Ok(()) => format!(
268 "reload failed, previous file restored and the service was not reloaded again: {e}"
269 ),
270 Err(x) => format!("reload failed AND restoring the previous file failed ({x}): {e}"),
271 });
272 }
273 }
274 Ok(Applied::Written)
275}
276
277#[cfg(test)]
278mod tests {
279 use super::*;
280
281 fn scratch(tag: &str) -> PathBuf {
282 let n = std::time::SystemTime::now()
283 .duration_since(std::time::UNIX_EPOCH)
284 .unwrap()
285 .as_nanos();
286 let d = std::env::temp_dir().join(format!("unv-apply-{tag}-{n}"));
287 std::fs::create_dir_all(&d).unwrap();
288 d
289 }
290
291 fn target(dir: &Path, validate: Option<&str>, reload: Option<&str>) -> Target {
292 Target {
293 id: "t".into(),
294 path: dir.join("app.conf"),
295 project: "p".into(),
296 exporter: "nginx".into(),
297 mode: "push".into(),
298 apply: true,
299 validate: validate.map(String::from),
300 reload: reload.map(String::from),
301 require_approval: false,
302 }
303 }
304
305 #[test]
306 fn a_hash_mismatch_writes_nothing() {
307 let d = scratch("hash");
308 let t = target(&d, None, None);
309 let e = apply(&t, b"new", &sha256_hex(b"other"), &d).unwrap_err();
310 assert!(e.contains("does not match"));
311 assert!(!t.path.exists());
312 }
313
314 #[test]
315 fn apply_is_refused_unless_the_nodes_own_config_allows_it() {
316 let d = scratch("noapply");
317 let mut t = target(&d, None, None);
318 t.apply = false;
319 assert!(apply(&t, b"x", &sha256_hex(b"x"), &d)
320 .unwrap_err()
321 .contains("not set to apply"));
322 t.apply = true;
323 t.mode = "pull".into();
324 assert!(apply(&t, b"x", &sha256_hex(b"x"), &d).is_err());
325 assert!(!t.path.exists());
326 }
327
328 #[test]
329 fn writes_the_file_and_keeps_the_previous_one() {
330 let d = scratch("write");
331 let t = target(&d, None, None);
332 std::fs::write(&t.path, "old").unwrap();
333 assert_eq!(
334 apply(&t, b"new", &sha256_hex(b"new"), &d).unwrap(),
335 Applied::Written
336 );
337 assert_eq!(std::fs::read_to_string(&t.path).unwrap(), "new");
338 let kept = kept_versions(&d, "t");
339 assert_eq!(kept.len(), 1);
340 assert_eq!(std::fs::read_to_string(&kept[0]).unwrap(), "old");
341 assert_eq!(
343 apply(&t, b"new", &sha256_hex(b"new"), &d).unwrap(),
344 Applied::Unchanged
345 );
346 assert_eq!(kept_versions(&d, "t").len(), 1);
347 }
348
349 #[test]
350 fn only_the_last_three_versions_are_kept() {
351 let d = scratch("keep");
352 let t = target(&d, None, None);
353 std::fs::write(&t.path, "v0").unwrap();
354 for i in 1..=6 {
355 let c = format!("v{i}");
356 apply(&t, c.as_bytes(), &sha256_hex(c.as_bytes()), &d).unwrap();
359 }
360 assert_eq!(kept_versions(&d, "t").len(), KEEP_VERSIONS);
361 }
362
363 #[cfg(unix)]
364 #[test]
365 fn a_failing_validate_restores_the_old_file_and_never_reloads() {
366 let d = scratch("validate");
367 let marker = d.join("reloaded");
368 let t = target(
369 &d,
370 Some("echo bad config >&2; exit 1"),
371 Some(&format!("touch {}", marker.display())),
372 );
373 std::fs::write(&t.path, "old").unwrap();
374 let e = apply(&t, b"new", &sha256_hex(b"new"), &d).unwrap_err();
375 assert!(e.contains("validate failed, previous file restored"), "{e}");
376 assert!(e.contains("bad config"), "{e}");
377 assert_eq!(std::fs::read_to_string(&t.path).unwrap(), "old");
378 assert!(!marker.exists(), "reload ran after a failed validate");
379 }
380
381 #[test]
382 fn a_failing_validate_on_a_new_file_removes_it() {
383 let d = scratch("validate-new");
384 let t = target(&d, Some("exit 3"), None);
385 assert!(apply(&t, b"new", &sha256_hex(b"new"), &d).is_err());
386 assert!(!t.path.exists());
387 }
388
389 #[test]
390 fn a_failing_reload_restores_the_old_file() {
391 let d = scratch("reload");
392 let t = target(&d, Some("exit 0"), Some("exit 2"));
393 std::fs::write(&t.path, "old").unwrap();
394 let e = apply(&t, b"new", &sha256_hex(b"new"), &d).unwrap_err();
395 assert!(e.contains("reload failed, previous file restored"), "{e}");
396 assert_eq!(std::fs::read_to_string(&t.path).unwrap(), "old");
397 }
398
399 #[cfg(unix)]
400 #[test]
401 fn validate_sees_the_new_file_in_place() {
402 let d = scratch("inplace");
403 let t = target(
404 &d,
405 Some(&format!("grep -q NEW {}", d.join("app.conf").display())),
406 None,
407 );
408 std::fs::write(&t.path, "old").unwrap();
409 assert!(apply(&t, b"NEW", &sha256_hex(b"NEW"), &d).is_ok());
410 }
411
412 #[cfg(unix)]
413 #[test]
414 fn a_runaway_command_is_killed_at_the_timeout() {
415 let t0 = Instant::now();
416 let e = run("sleep 30", "", Duration::from_millis(300)).unwrap_err();
417 assert!(e.contains("did not finish"), "{e}");
418 assert!(
419 t0.elapsed() < Duration::from_secs(5),
420 "waited for the command"
421 );
422 }
423
424 #[cfg(unix)]
425 #[test]
426 fn an_existing_files_mode_is_kept_and_a_new_one_is_private() {
427 use std::os::unix::fs::PermissionsExt;
428 let d = scratch("mode");
429 let t = target(&d, None, None);
430 std::fs::write(&t.path, "old").unwrap();
431 std::fs::set_permissions(&t.path, std::fs::Permissions::from_mode(0o640)).unwrap();
432 apply(&t, b"new", &sha256_hex(b"new"), &d).unwrap();
433 assert_eq!(
434 std::fs::metadata(&t.path).unwrap().permissions().mode() & 0o777,
435 0o640
436 );
437
438 let t2 = Target {
439 id: "u".into(),
440 path: d.join("fresh.conf"),
441 ..t
442 };
443 apply(&t2, b"x", &sha256_hex(b"x"), &d).unwrap();
444 assert_eq!(
445 std::fs::metadata(&t2.path).unwrap().permissions().mode() & 0o777,
446 0o600
447 );
448 let kept = kept_versions(&d, "t");
450 assert_eq!(
451 std::fs::metadata(&kept[0]).unwrap().permissions().mode() & 0o777,
452 0o600
453 );
454 }
455
456 #[test]
457 fn scrub_elides_a_line_of_the_file_and_keeps_the_rest() {
458 let content = "server {\n secret_token abcdef123456;\n}\n";
459 let out = scrub(
460 "nginx: [emerg] unknown directive in conf:2\nsecret_token abcdef123456;\n",
461 content,
462 );
463 assert!(!out.contains("abcdef123456"));
464 assert!(out.contains("unknown directive"));
465 assert!(out.contains("elided"));
466 }
467
468 #[test]
469 fn a_missing_file_has_no_hash_and_a_present_one_does() {
470 let d = scratch("hashfile");
471 assert_eq!(hash_file(&d.join("nope")).unwrap(), None);
472 std::fs::write(d.join("f"), "abc").unwrap();
473 assert_eq!(
474 hash_file(&d.join("f")).unwrap().unwrap(),
475 sha256_hex(b"abc")
476 );
477 }
478}