1use anyhow::Error;
8use argh::FromArgs;
9use diagnostics_reader::{ArchiveReader, InspectArchiveReader};
10use fidl::endpoints::Proxy;
11use fidl_fuchsia_diagnostics as fdiagnostics;
12use fidl_fuchsia_diagnostics_persistence as fpersistence;
13use fidl_fuchsia_io as fio;
14use fidl_fuchsia_power_battery as fbattery;
15use fidl_fuchsia_update as fupdate;
16use fuchsia_async as fasync;
17use fuchsia_component::client::{connect_to_protocol, connect_to_protocol_at_path};
18use fuchsia_component::server::ServiceFs;
19use fuchsia_inspect::health::Reporter;
20use futures::lock::Mutex;
21use futures::{FutureExt, StreamExt, TryStreamExt};
22use log::*;
23use persistence_build_config::Config;
24use serde::{Deserialize, Serialize};
25use std::path::{Path, PathBuf};
26use std::sync::Arc;
27use zx::{BootInstant, MonotonicDuration, MonotonicInstant};
28
29pub const PROGRAM_NAME: &str = "persistence";
31
32const BATTERY_TRIGGER_HYSTERESIS_PERCENT: f32 = 5.0;
34
35#[derive(FromArgs, Debug, PartialEq)]
37#[argh(subcommand, name = "persistence")]
38pub struct CommandLine {}
39
40#[derive(Serialize, Deserialize, Debug, Clone, PartialEq)]
41pub struct Metadata {
42 pub monotonic_timestamp: i64,
43 pub boot_timestamp: i64,
44 pub utc_timestamp_ns: Option<i64>,
45}
46
47enum IncomingService {
48 PreviousBootDataProvider(fpersistence::PreviousBootDataProviderRequestStream),
49}
50
51pub async fn main(_args: CommandLine) -> Result<(), Error> {
52 info!("Persistence starting up");
53 let scope = fasync::Scope::new();
54
55 let _inspect_controller = inspect_runtime::publish(
56 fuchsia_inspect::component::inspector(),
57 inspect_runtime::PublishOptions::default().custom_scope(scope.clone()),
58 );
59
60 fuchsia_inspect::component::health().set_starting_up();
61 let config = Config::take_from_startup_handle();
62
63 let cache_dir = Path::new("/cache");
64 if let Err(e) = rotate_active_to_previous_boot(cache_dir) {
65 warn!(e:?; "Error rotating active data to previous boot");
66 }
67
68 let (ready_tx, ready_rx) = futures::channel::oneshot::channel::<()>();
69 let ready = ready_rx.shared();
70
71 let mut fs = ServiceFs::new();
72 fs.dir("svc").add_fidl_service(IncomingService::PreviousBootDataProvider);
73 fs.take_and_serve_directory_handle()?;
74
75 let cache_dir_buf = cache_dir.to_path_buf();
76 let ready_for_handler = ready.clone();
77 scope.spawn(async move {
78 fs.for_each_concurrent(None, |IncomingService::PreviousBootDataProvider(stream)| {
79 let cache_dir = cache_dir_buf.clone();
80 let ready = ready_for_handler.clone();
81 async move {
82 if let Err(e) = handle_previous_boot_data_provider(stream, &cache_dir, ready).await
83 {
84 warn!(e:?; "Error handling PreviousBootDataProvider request stream");
85 }
86 }
87 })
88 .await;
89 });
90
91 if config.skip_update_check {
92 info!("Skipping the update check, publishing previous boot data");
93 let _ = ready_tx.send(());
94 } else {
95 scope.spawn(async move {
96 if let Err(e) = wait_for_update().await {
97 warn!(e:?; "Will not publish previous boot data");
98 } else {
99 let _ = ready_tx.send(());
100 }
101 });
102 }
103
104 let proxy = connect_to_protocol_at_path::<fdiagnostics::ArchiveAccessorMarker>(
105 "/svc/fuchsia.diagnostics.ArchiveAccessor.previous_boot",
106 )?;
107
108 let mut reader = ArchiveReader::inspect();
109 reader.with_archive(proxy);
110 let collector = Arc::new(Mutex::new(SnapshotCollector::new(reader, cache_dir.to_path_buf())));
111
112 let period = MonotonicDuration::from_seconds(config.persistence_period_seconds);
113 let collector_periodic = collector.clone();
114 scope.spawn(async move {
115 if let Err(e) = collector_periodic.lock().await.collect_active_snapshot().await {
116 error!(e:?; "Error collecting initial active inspect snapshot");
117 }
118 let mut interval = fasync::Interval::new(period);
119 while let Some(()) = interval.next().await {
120 if let Err(e) = collector_periodic.lock().await.collect_active_snapshot().await {
121 error!(e:?; "Error collecting active inspect snapshot");
122 }
123 }
124 });
125
126 if let Ok(battery_manager) = connect_to_protocol::<fbattery::BatteryManagerMarker>() {
127 let threshold = config.low_battery_threshold_percent as f32;
128 let collector_battery = collector.clone();
129 scope.spawn(async move {
130 if let Err(e) =
131 listen_for_low_battery(battery_manager, threshold, collector_battery).await
132 {
133 warn!(e:?; "BatteryManager watcher task terminated");
134 }
135 });
136 }
137
138 fuchsia_inspect::component::health().set_ok();
139 scope.await;
140
141 Ok(())
142}
143
144async fn listen_for_low_battery(
145 battery_manager: fbattery::BatteryManagerProxy,
146 threshold_percent: f32,
147 collector: Arc<Mutex<SnapshotCollector>>,
148) -> Result<(), Error> {
149 let (watcher_client, mut request_stream) =
150 fidl::endpoints::create_request_stream::<fbattery::BatteryInfoWatcherMarker>();
151 battery_manager.watch(watcher_client)?;
152
153 let mut triggered = false;
154
155 while let Some(request) = request_stream.try_next().await? {
156 match request {
157 fbattery::BatteryInfoWatcherRequest::OnChangeBatteryInfo {
158 info,
159 wake_lease: _,
160 responder,
161 } => {
162 let level_percent = info.level_percent.unwrap_or(100.0);
163 let level_status = info.level_status.unwrap_or(fbattery::LevelStatus::Unknown);
164
165 let is_low = level_percent <= threshold_percent
166 || level_status == fbattery::LevelStatus::Low
167 || level_status == fbattery::LevelStatus::Critical;
168
169 if is_low {
170 if !triggered {
171 info!(
172 level_percent,
173 level_status:?;
174 "Low battery threshold reached, collecting active snapshot"
175 );
176 if let Err(e) = collector.lock().await.collect_active_snapshot().await {
177 error!(e:?; "Error collecting active snapshot on low battery");
178 }
179 triggered = true;
180 }
181 } else if level_percent > threshold_percent + BATTERY_TRIGGER_HYSTERESIS_PERCENT {
182 triggered = false;
183 }
184
185 responder.send()?;
186 }
187 }
188 }
189 Ok(())
190}
191
192fn rotate_active_to_previous_boot(cache_dir: &Path) -> Result<(), Error> {
193 let active_dir = cache_dir.join("active");
194 let previous_boot_dir = cache_dir.join("previous_boot");
195
196 std::fs::create_dir_all(&active_dir)?;
197 std::fs::create_dir_all(&previous_boot_dir)?;
198
199 let active_file = active_dir.join("active.json");
200 let active_meta = active_dir.join("metadata.json");
201
202 info!(
203 active_exists = active_file.exists(),
204 meta_exists = active_meta.exists();
205 "Rotating active data to previous boot"
206 );
207
208 for entry in std::fs::read_dir(&previous_boot_dir)? {
210 let entry = entry?;
211 let path = entry.path();
212 if path.is_file() {
213 let _ = std::fs::remove_file(path);
214 }
215 }
216
217 if active_file.exists() {
218 std::fs::rename(&active_file, previous_boot_dir.join("active.json"))?;
219 }
220 if active_meta.exists() {
221 std::fs::rename(&active_meta, previous_boot_dir.join("metadata.json"))?;
222 }
223
224 for entry in std::fs::read_dir(&active_dir)? {
226 let entry = entry?;
227 let path = entry.path();
228 if path.is_file() {
229 let _ = std::fs::remove_file(path);
230 }
231 }
232
233 Ok(())
234}
235
236fn maybe_get_utc_timestamp_from_clock<R: zx::Timeline, O: zx::Timeline>(
237 clock: Option<&zx::Clock<R, O>>,
238) -> Option<i64> {
239 let clock = clock?;
240 clock.read().ok().map(|time| time.into_nanos())
241}
242
243pub fn maybe_get_utc_timestamp() -> Option<i64> {
244 let clock = fuchsia_runtime::duplicate_utc_clock_handle(zx::Rights::SAME_RIGHTS).ok();
245 maybe_get_utc_timestamp_from_clock(clock.as_ref())
246}
247
248struct SnapshotCollector {
249 reader: InspectArchiveReader,
250 cache_dir: PathBuf,
251}
252
253impl SnapshotCollector {
254 fn new(reader: InspectArchiveReader, cache_dir: PathBuf) -> Self {
255 Self { reader, cache_dir }
256 }
257
258 async fn collect_active_snapshot(&mut self) -> Result<(), Error> {
259 info!("Collecting active snapshot...");
260 let inspect_data = match self.reader.snapshot().await {
261 Ok(v) => v,
262 Err(e) => {
263 error!(e:?; "ArchiveReader snapshot failed");
264 return Err(e.into());
265 }
266 };
267 let active_dir = self.cache_dir.join("active");
268 std::fs::create_dir_all(&active_dir)?;
269
270 let active_json_data = serde_json::to_vec(&inspect_data)?;
271 let tmp_file = active_dir.join("active.json.tmp");
272 let final_file = active_dir.join("active.json");
273 std::fs::write(&tmp_file, active_json_data)?;
274 std::fs::rename(&tmp_file, &final_file)?;
275
276 let metadata = Metadata {
277 monotonic_timestamp: MonotonicInstant::get().into_nanos(),
278 boot_timestamp: BootInstant::get().into_nanos(),
279 utc_timestamp_ns: maybe_get_utc_timestamp(),
280 };
281
282 let meta_json_data = serde_json::to_vec(&metadata)?;
283 let tmp_meta = active_dir.join("metadata.json.tmp");
284 let final_meta = active_dir.join("metadata.json");
285 std::fs::write(&tmp_meta, meta_json_data)?;
286 std::fs::rename(&tmp_meta, &final_meta)?;
287
288 Ok(())
289 }
290}
291
292async fn handle_previous_boot_data_provider(
293 mut stream: fpersistence::PreviousBootDataProviderRequestStream,
294 cache_dir: &Path,
295 ready: futures::future::Shared<futures::channel::oneshot::Receiver<()>>,
296) -> Result<(), Error> {
297 while let Some(request) = stream.try_next().await? {
298 match request {
299 fpersistence::PreviousBootDataProviderRequest::WatchPreviousBootData {
300 options: _,
301 responder,
302 } => {
303 info!("Received WatchPreviousBootData request, awaiting ready signal...");
304 if ready.clone().await.is_err() {
305 warn!("Update check failed or cancelled; will not serve previous boot data");
306 return Ok(());
307 }
308 info!("Ready signal received in WatchPreviousBootData handler");
309 let previous_boot_dir = cache_dir.join("previous_boot");
310 let active_file = previous_boot_dir.join("active.json");
311 let meta_file = previous_boot_dir.join("metadata.json");
312
313 if !active_file.exists() || !meta_file.exists() {
314 responder.send(fpersistence::PreviousBootData::default())?;
315 continue;
316 }
317
318 let meta_bytes = std::fs::read(&meta_file)?;
319 let metadata: Metadata = serde_json::from_slice(&meta_bytes)?;
320
321 let active_path_str =
322 active_file.to_str().ok_or_else(|| anyhow::anyhow!("Invalid path"))?;
323 let file_proxy =
324 fuchsia_fs::file::open_in_namespace(active_path_str, fio::PERM_READABLE)?;
325 let client_end = file_proxy
326 .into_client_end()
327 .map_err(|_| anyhow::anyhow!("Failed to convert to client end"))?;
328
329 let data = fpersistence::PreviousBootData {
330 monotonic_timestamp: Some(MonotonicInstant::from_nanos(
331 metadata.monotonic_timestamp,
332 )),
333 boot_timestamp: Some(BootInstant::from_nanos(metadata.boot_timestamp)),
334 utc_timestamp_ns: metadata.utc_timestamp_ns,
335 inspect: Some(client_end),
336 ..Default::default()
337 };
338
339 responder.send(data)?;
340 }
341 fpersistence::PreviousBootDataProviderRequest::_UnknownMethod { .. } => {}
342 }
343 }
344 Ok(())
345}
346
347async fn wait_for_update() -> Result<(), Error> {
348 info!("Waiting for post-boot update check...");
349 let (notifier_client, mut notifier_request_stream) =
350 fidl::endpoints::create_request_stream::<fupdate::NotifierMarker>();
351 let proxy = connect_to_protocol::<fupdate::ListenerMarker>()?;
352 proxy.notify_on_first_update_check(fupdate::ListenerNotifyOnFirstUpdateCheckRequest {
353 notifier: Some(notifier_client),
354 ..Default::default()
355 })?;
356
357 match notifier_request_stream.try_next().await {
358 Ok(Some(fupdate::NotifierRequest::Notify { control_handle: _ })) => {}
359 Ok(None) => {
360 return Err(anyhow::anyhow!("Did not receive update notification; not publishing"));
361 }
362 Err(e) => {
363 return Err(anyhow::anyhow!(
364 "Error waiting for update notification; not publishing: {e}"
365 ));
366 }
367 }
368
369 info!("...Update check has completed; publishing previous boot data");
371 Ok(())
372}
373
374#[cfg(test)]
375mod tests {
376 use super::*;
377 use tempfile::TempDir;
378
379 #[test]
380 fn test_metadata_serde() {
381 let meta = Metadata {
382 monotonic_timestamp: 1000,
383 boot_timestamp: 2000,
384 utc_timestamp_ns: Some(3000),
385 };
386
387 let serialized = serde_json::to_vec(&meta).unwrap();
388 let deserialized: Metadata = serde_json::from_slice(&serialized).unwrap();
389 assert_eq!(meta, deserialized);
390 }
391
392 #[test]
393 fn test_metadata_serde_without_utc() {
394 let meta =
395 Metadata { monotonic_timestamp: 1000, boot_timestamp: 2000, utc_timestamp_ns: None };
396
397 let serialized = serde_json::to_vec(&meta).unwrap();
398 let deserialized: Metadata = serde_json::from_slice(&serialized).unwrap();
399 assert_eq!(meta, deserialized);
400 }
401
402 #[test]
403 fn test_maybe_get_utc_timestamp_when_utc_not_available() {
404 assert_eq!(
406 maybe_get_utc_timestamp_from_clock::<zx::MonotonicTimeline, zx::SyntheticTimeline>(
407 None
408 ),
409 None
410 );
411
412 let clock = zx::SyntheticClock::create(zx::ClockOpts::AUTO_START, None)
414 .expect("failed to create clock");
415 let no_read_clock = clock.duplicate_handle(zx::Rights::NONE).expect("duplicate handle");
416 assert_eq!(maybe_get_utc_timestamp_from_clock(Some(&no_read_clock)), None);
417
418 assert!(maybe_get_utc_timestamp_from_clock(Some(&clock)).is_some());
420 }
421
422 #[test]
423 fn test_rotation() {
424 let _ = std::fs::create_dir_all("/cache");
425 let temp_dir = TempDir::new_in("/cache").unwrap();
426 let cache_dir = temp_dir.path();
427 let active_dir = cache_dir.join("active");
428 let previous_boot_dir = cache_dir.join("previous_boot");
429
430 std::fs::create_dir_all(&active_dir).unwrap();
431 std::fs::write(active_dir.join("active.json"), b"active_data").unwrap();
432 std::fs::write(active_dir.join("metadata.json"), b"meta_data").unwrap();
433
434 rotate_active_to_previous_boot(cache_dir).unwrap();
435
436 assert!(!active_dir.join("active.json").exists());
437 assert!(!active_dir.join("metadata.json").exists());
438 assert_eq!(std::fs::read(previous_boot_dir.join("active.json")).unwrap(), b"active_data");
439 assert_eq!(std::fs::read(previous_boot_dir.join("metadata.json")).unwrap(), b"meta_data");
440 }
441
442 #[test]
443 fn test_rotation_cleans_non_clean_active_dir() {
444 let _ = std::fs::create_dir_all("/cache");
445 let temp_dir = TempDir::new_in("/cache").unwrap();
446 let cache_dir = temp_dir.path();
447 let active_dir = cache_dir.join("active");
448 let previous_boot_dir = cache_dir.join("previous_boot");
449
450 std::fs::create_dir_all(&active_dir).unwrap();
451 std::fs::write(active_dir.join("active.json"), b"active_data").unwrap();
452 std::fs::write(active_dir.join("metadata.json"), b"meta_data").unwrap();
453 std::fs::write(active_dir.join("active.json.tmp"), b"tmp_active_data").unwrap();
455 std::fs::write(active_dir.join("metadata.json.tmp"), b"tmp_meta_data").unwrap();
456 std::fs::write(active_dir.join("stray.txt"), b"stray_data").unwrap();
457
458 rotate_active_to_previous_boot(cache_dir).unwrap();
459
460 assert_eq!(std::fs::read(previous_boot_dir.join("active.json")).unwrap(), b"active_data");
462 assert_eq!(std::fs::read(previous_boot_dir.join("metadata.json")).unwrap(), b"meta_data");
463 assert!(!previous_boot_dir.join("active.json.tmp").exists());
465 assert!(!previous_boot_dir.join("metadata.json.tmp").exists());
466 assert!(!previous_boot_dir.join("stray.txt").exists());
467
468 assert_eq!(std::fs::read_dir(&active_dir).unwrap().count(), 0);
470 }
471
472 #[test]
473 fn test_rotation_clears_stale_previous_boot_when_active_empty() {
474 let _ = std::fs::create_dir_all("/cache");
475 let temp_dir = TempDir::new_in("/cache").unwrap();
476 let cache_dir = temp_dir.path();
477 let active_dir = cache_dir.join("active");
478 let previous_boot_dir = cache_dir.join("previous_boot");
479
480 std::fs::create_dir_all(&active_dir).unwrap();
481 std::fs::create_dir_all(&previous_boot_dir).unwrap();
482
483 std::fs::write(previous_boot_dir.join("active.json"), b"stale_active_data").unwrap();
485 std::fs::write(previous_boot_dir.join("metadata.json"), b"stale_meta_data").unwrap();
486
487 rotate_active_to_previous_boot(cache_dir).unwrap();
489
490 assert!(!previous_boot_dir.join("active.json").exists());
492 assert!(!previous_boot_dir.join("metadata.json").exists());
493 }
494
495 #[fasync::run_singlethreaded(test)]
496 async fn test_handle_previous_boot_data_provider_missing_files() {
497 let _ = std::fs::create_dir_all("/cache");
498 let temp_dir = TempDir::new_in("/cache").unwrap();
499 let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<
500 fpersistence::PreviousBootDataProviderMarker,
501 >();
502
503 let (ready_tx, ready_rx) = futures::channel::oneshot::channel();
504 let _ = ready_tx.send(());
505 let ready = ready_rx.shared();
506
507 let scope = fasync::Scope::new();
508 let cache_dir = temp_dir.path().to_path_buf();
509 scope.spawn(async move {
510 let _ = handle_previous_boot_data_provider(stream, &cache_dir, ready).await;
511 });
512
513 let data = proxy.watch_previous_boot_data(&Default::default()).await.unwrap();
514 assert!(data.inspect.is_none());
515 assert!(data.monotonic_timestamp.is_none());
516 assert!(data.boot_timestamp.is_none());
517 assert!(data.utc_timestamp_ns.is_none());
518 assert_eq!(data, fpersistence::PreviousBootData::default());
519 }
520
521 #[fasync::run_singlethreaded(test)]
522 async fn test_handle_previous_boot_data_provider_valid_data() {
523 let _ = std::fs::create_dir_all("/cache");
524 let temp_dir = TempDir::new_in("/cache").unwrap();
525 let previous_boot_dir = temp_dir.path().join("previous_boot");
526 std::fs::create_dir_all(&previous_boot_dir).unwrap();
527
528 let meta = Metadata {
529 monotonic_timestamp: 12345,
530 boot_timestamp: 67890,
531 utc_timestamp_ns: Some(100000),
532 };
533 let meta_bytes = serde_json::to_vec(&meta).unwrap();
534 std::fs::write(previous_boot_dir.join("metadata.json"), meta_bytes).unwrap();
535 std::fs::write(previous_boot_dir.join("active.json"), b"valid_inspect_json").unwrap();
536
537 let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<
538 fpersistence::PreviousBootDataProviderMarker,
539 >();
540
541 let (ready_tx, ready_rx) = futures::channel::oneshot::channel();
542 let _ = ready_tx.send(());
543 let ready = ready_rx.shared();
544
545 let scope = fasync::Scope::new();
546 let cache_dir = temp_dir.path().to_path_buf();
547 scope.spawn(async move {
548 let _ = handle_previous_boot_data_provider(stream, &cache_dir, ready).await;
549 });
550
551 let data = proxy.watch_previous_boot_data(&Default::default()).await.unwrap();
552 assert_eq!(data.monotonic_timestamp, Some(MonotonicInstant::from_nanos(12345)));
553 assert_eq!(data.boot_timestamp, Some(BootInstant::from_nanos(67890)));
554 assert_eq!(data.utc_timestamp_ns, Some(100000));
555 assert!(data.inspect.is_some());
556 }
557
558 #[fasync::run_singlethreaded(test)]
559 async fn test_handle_previous_boot_data_provider_utc_not_available() {
560 let _ = std::fs::create_dir_all("/cache");
561 let temp_dir = TempDir::new_in("/cache").unwrap();
562 let previous_boot_dir = temp_dir.path().join("previous_boot");
563 std::fs::create_dir_all(&previous_boot_dir).unwrap();
564
565 let meta =
566 Metadata { monotonic_timestamp: 12345, boot_timestamp: 67890, utc_timestamp_ns: None };
567 let meta_bytes = serde_json::to_vec(&meta).unwrap();
568 std::fs::write(previous_boot_dir.join("metadata.json"), meta_bytes).unwrap();
569 std::fs::write(previous_boot_dir.join("active.json"), b"valid_inspect_json").unwrap();
570
571 let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<
572 fpersistence::PreviousBootDataProviderMarker,
573 >();
574
575 let (ready_tx, ready_rx) = futures::channel::oneshot::channel();
576 let _ = ready_tx.send(());
577 let ready = ready_rx.shared();
578
579 let scope = fasync::Scope::new();
580 let cache_dir = temp_dir.path().to_path_buf();
581 scope.spawn(async move {
582 let _ = handle_previous_boot_data_provider(stream, &cache_dir, ready).await;
583 });
584
585 let data = proxy.watch_previous_boot_data(&Default::default()).await.unwrap();
586 assert_eq!(data.monotonic_timestamp, Some(MonotonicInstant::from_nanos(12345)));
587 assert_eq!(data.boot_timestamp, Some(BootInstant::from_nanos(67890)));
588 assert_eq!(data.utc_timestamp_ns, None);
589 assert!(data.inspect.is_some());
590 }
591
592 #[fasync::run_singlethreaded(test)]
593 async fn test_wait_for_update_no_listener() {
594 let res = wait_for_update().await;
595 assert!(res.is_err());
596 }
597
598 #[fasync::run_singlethreaded(test)]
599 async fn test_handle_previous_boot_data_provider_canceled_ready() {
600 let _ = std::fs::create_dir_all("/cache");
601 let temp_dir = TempDir::new_in("/cache").unwrap();
602 let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<
603 fpersistence::PreviousBootDataProviderMarker,
604 >();
605
606 let (ready_tx, ready_rx) = futures::channel::oneshot::channel::<()>();
607 drop(ready_tx);
608 let ready = ready_rx.shared();
609
610 let scope = fasync::Scope::new();
611 let cache_dir = temp_dir.path().to_path_buf();
612 scope.spawn(async move {
613 let _ = handle_previous_boot_data_provider(stream, &cache_dir, ready).await;
614 });
615
616 let res = proxy.watch_previous_boot_data(&Default::default()).await;
617 assert!(res.is_err());
618 }
619}