fake_archive_accessor/
fake_archive_accessor.rs1mod archivist_accessor;
6mod archivist_server;
7
8use {
9 anyhow::{
10 Error,
12 bail,
13 },
14 archivist_accessor::ArchiveAccessor,
15 async_trait::async_trait,
16 fidl_fuchsia_diagnostics as diagnostics, fuchsia_async as fasync,
17 fuchsia_sync::Mutex,
18 futures::StreamExt,
19 log::*,
20 std::collections::BTreeSet,
21 std::sync::{
22 Arc,
23 atomic::{AtomicUsize, Ordering},
24 },
25};
26
27#[async_trait]
30pub trait EventSignaler: Send + Sync {
31 async fn signal_fetch(&self);
33 async fn signal_done(&self);
35 async fn signal_error(&self, error: &str);
37}
38
39pub struct FakeArchiveAccessor {
43 event_signaler: Option<Box<dyn EventSignaler>>,
44 inspect_data: Vec<String>,
45 next_data: AtomicUsize,
46 selectors_requested: Mutex<Vec<BTreeSet<String>>>,
48}
49
50impl FakeArchiveAccessor {
51 pub fn new(
56 inspect_data: &[String],
57 event_signaler: Option<Box<dyn EventSignaler>>,
58 ) -> Arc<FakeArchiveAccessor> {
59 Arc::new(FakeArchiveAccessor {
60 inspect_data: inspect_data.to_vec(),
61 event_signaler,
62 next_data: AtomicUsize::new(0),
63 selectors_requested: Mutex::new(vec![]),
64 })
65 }
66
67 async fn handle_fidl_request(
70 &self,
71 request: diagnostics::ArchiveAccessorRequest,
72 ) -> Result<(), Error> {
73 let (stream_parameters, result_stream) = match request {
74 diagnostics::ArchiveAccessorRequest::StreamDiagnostics {
75 stream_parameters,
76 result_stream,
77 control_handle: _,
78 } => (stream_parameters, result_stream),
79 diagnostics::ArchiveAccessorRequest::WaitForReady { .. }
80 | diagnostics::ArchiveAccessorRequest::StreamDiagnosticsToSocket { .. }
81 | diagnostics::ArchiveAccessorRequest::_UnknownMethod { .. } => {
82 unreachable!("Unexpected method call");
83 }
84 };
85 let selectors = ArchiveAccessor::validate_stream_request(stream_parameters)?;
86 self.selectors_requested.lock().push(
87 selectors
88 .into_iter()
89 .map(|s| {
90 selectors::selector_to_string(
91 &s,
92 selectors::SelectorDisplayOptions::never_wrap_in_quotes(),
93 )
94 })
95 .collect::<Result<BTreeSet<_>, Error>>()?,
96 );
97 if let Some(s) = self.event_signaler.as_ref() {
98 s.signal_fetch().await;
99 }
100 let data_index = self.next_data.fetch_add(1, Ordering::Relaxed);
101 if data_index >= self.inspect_data.len() {
102 if let Some(s) = self.event_signaler.as_ref() {
105 s.signal_done().await;
106 }
107 } else if let Err(problem) =
108 ArchiveAccessor::send(result_stream, &self.inspect_data[data_index]).await
109 {
110 if let Some(s) = self.event_signaler.as_ref() {
111 s.signal_done().await;
112 s.signal_error(&format!("{problem}")).await;
113 }
114 error!("Problem in request: {}", problem);
115 return Err(problem);
116 }
117 Ok(())
118 }
119
120 pub async fn serve_stream(
121 &self,
122 mut request_stream: diagnostics::ArchiveAccessorRequestStream,
123 ) -> Result<(), Error> {
124 loop {
125 match request_stream.next().await {
126 Some(Ok(request)) => self.handle_fidl_request(request).await?,
127 Some(Err(e)) => {
128 if let Some(s) = self.event_signaler.as_ref() {
129 s.signal_done().await;
130 }
131 bail!("{}", e);
132 }
133 None => break,
134 }
135 }
136 Ok(())
137 }
138
139 pub fn serve_async(self: Arc<Self>, stream: diagnostics::ArchiveAccessorRequestStream) {
140 fasync::Task::spawn(async move {
141 let result = self.serve_stream(stream).await;
142 if let Err(e) = result {
143 error!("Error while serving ArchiveAccessor: {:?}", e);
144 }
145 })
146 .detach();
147 }
148
149 pub fn get_selectors_requested(&self) -> Vec<BTreeSet<String>> {
150 self.selectors_requested.lock().clone()
151 }
152}