1use crate::accessor::{ArchiveAccessorServer, ArchiveAccessorTranslator, ArchiveAccessorWriter};
5use crate::configs;
6use crate::diagnostics::AccessorStats;
7use crate::error::Error;
8use crate::pipeline::StaticHierarchyAllowlist;
9use fidl::endpoints::{DiscoverableProtocolMarker, ProtocolMarker, ServerEnd};
10use fidl_fuchsia_component_sandbox as fsandbox;
11use fidl_fuchsia_diagnostics as fdiagnostics;
12use fidl_fuchsia_diagnostics_host as fdiagnostics_host;
13use fidl_fuchsia_io as fio;
14use fuchsia_async as fasync;
15use fuchsia_fs::directory;
16use fuchsia_inspect as inspect;
17use fuchsia_sync::RwLock;
18use futures::TryStreamExt;
19use log::{debug, warn};
20use moniker::ExtendedMoniker;
21use std::borrow::Cow;
22use std::ops::Deref;
23use std::path::{Path, PathBuf};
24use std::sync::{Arc, Weak};
25
26const ALL_PIPELINE_NAME: &str = "all";
27
28struct PipelineParameters {
29 has_config: bool,
30 name: Cow<'static, str>,
31 empty_behavior: configs::EmptyBehavior,
32}
33
34#[derive(Copy, Clone)]
35pub struct PipelinesDictionaryId(u64);
36
37impl Deref for PipelinesDictionaryId {
38 type Target = u64;
39
40 fn deref(&self) -> &Self::Target {
41 &self.0
42 }
43}
44
45pub struct Pipeline {
50 name: Cow<'static, str>,
52
53 _pipeline_node: Option<inspect::Node>,
55
56 stats: AccessorStats,
58
59 static_allowlist: RwLock<StaticHierarchyAllowlist>,
61}
62
63impl Pipeline {
64 fn new(
65 parameters: PipelineParameters,
66 pipelines_path: &Path,
67 parent_node: &inspect::Node,
68 accessor_stats_node: &inspect::Node,
69 ) -> Self {
70 let mut _pipeline_node = None;
71 let path = format!("{}/{}", pipelines_path.display(), parameters.name);
72 let mut static_selectors = None;
73 if parameters.has_config {
74 let node = parent_node.create_child(parameters.name.as_ref());
75 let mut config =
76 configs::PipelineConfig::from_directory(path, parameters.empty_behavior);
77 config.record_to_inspect(&node);
78 _pipeline_node = Some(node);
79 if !config.disable_filtering {
80 static_selectors = config.take_inspect_selectors();
81 }
82 }
83 let stats = AccessorStats::new(accessor_stats_node.create_child(parameters.name.as_ref()));
84 Pipeline {
85 _pipeline_node,
86 stats,
87 name: parameters.name,
88 static_allowlist: RwLock::new(StaticHierarchyAllowlist::new(static_selectors)),
89 }
90 }
91
92 fn protocol_name(&self) -> Cow<'_, str> {
93 self.protocol_name_inner::<fdiagnostics::ArchiveAccessorMarker>()
94 }
95
96 fn host_protocol_name(&self) -> Cow<'_, str> {
97 self.protocol_name_inner::<fdiagnostics_host::ArchiveAccessorMarker>()
98 }
99
100 fn protocol_name_inner<P: DiscoverableProtocolMarker>(&self) -> Cow<'_, str> {
101 if self.name.as_ref() == ALL_PIPELINE_NAME {
102 Cow::Borrowed(P::PROTOCOL_NAME)
103 } else {
104 Cow::Owned(format!("{}.{}", P::PROTOCOL_NAME, self.name))
105 }
106 }
107
108 pub fn accessor_stats(&self) -> &AccessorStats {
109 &self.stats
110 }
111
112 pub fn remove_component(&self, moniker: &ExtendedMoniker) {
113 self.static_allowlist.write().remove_component(moniker);
114 }
115
116 pub fn add_component(&self, moniker: &ExtendedMoniker) -> Result<(), Error> {
117 self.static_allowlist.write().add_component(moniker.clone())
118 }
119
120 pub fn static_hierarchy_allowlist(&self) -> StaticHierarchyAllowlist {
121 self.static_allowlist.read().clone()
126 }
127}
128
129#[cfg(test)]
130impl Pipeline {
131 pub fn for_test(static_selectors: Option<Vec<fdiagnostics::Selector>>) -> Self {
132 Pipeline {
133 _pipeline_node: None,
134 name: Cow::Borrowed("test"),
135 stats: AccessorStats::new(Default::default()),
136 static_allowlist: RwLock::new(StaticHierarchyAllowlist::new(static_selectors)),
137 }
138 }
139}
140
141pub struct PipelineManager {
142 pipelines: Vec<Arc<Pipeline>>,
143 _pipelines_node: inspect::Node,
144 _accessor_stats_node: inspect::Node,
145 scope: Option<fasync::Scope>,
146}
147
148impl PipelineManager {
149 pub async fn new(
150 pipelines_path: PathBuf,
151 pipelines_node: inspect::Node,
152 accessor_stats_node: inspect::Node,
153 scope: fasync::Scope,
154 ) -> Self {
155 let mut pipelines = vec![];
156 if let Ok(dir) =
157 directory::open_in_namespace(pipelines_path.to_str().unwrap(), fio::PERM_READABLE)
158 {
159 for entry in directory::readdir(&dir).await.expect("read dir") {
160 if !matches!(entry.kind, directory::DirentKind::Directory) {
161 continue;
162 }
163 let empty_behavior = if entry.name == "feedback" || entry.name == "previous_boot" {
164 configs::EmptyBehavior::DoNotFilter
165 } else {
166 configs::EmptyBehavior::Disable
167 };
168 let parameters = PipelineParameters {
169 has_config: true,
170 name: Cow::Owned(entry.name),
171 empty_behavior,
172 };
173 pipelines.push(Arc::new(Pipeline::new(
174 parameters,
175 &pipelines_path,
176 &pipelines_node,
177 &accessor_stats_node,
178 )));
179 }
180 }
181 pipelines.push(Arc::new(Pipeline::new(
182 PipelineParameters {
183 has_config: false,
184 name: Cow::Borrowed(ALL_PIPELINE_NAME),
185 empty_behavior: configs::EmptyBehavior::Disable,
186 },
187 &pipelines_path,
188 &pipelines_node,
189 &accessor_stats_node,
190 )));
191 Self {
192 pipelines,
193 _pipelines_node: pipelines_node,
194 _accessor_stats_node: accessor_stats_node,
195 scope: Some(scope),
196 }
197 }
198
199 pub fn weak_pipelines(&self) -> Vec<Weak<Pipeline>> {
200 self.pipelines.iter().map(Arc::downgrade).collect::<Vec<_>>()
201 }
202
203 pub async fn cancel(&mut self) {
204 if let Some(scope) = self.scope.take() {
205 scope.cancel().await;
206 }
207 }
208
209 pub async fn serve_pipelines(
210 &self,
211 accessor_server: Arc<ArchiveAccessorServer>,
212 id_gen: &sandbox::CapabilityIdGenerator,
213 capability_store: &mut fsandbox::CapabilityStoreProxy,
214 ) -> PipelinesDictionaryId {
215 let accessors_dict_id = id_gen.next();
216 capability_store.dictionary_create(accessors_dict_id).await.unwrap().unwrap();
217 debug!("Will serve {} pipelines", self.pipelines.len());
218 for pipeline in &self.pipelines {
219 debug!("Installing spawning receivers for {}", pipeline.name);
220 let accessor_pipeline = Arc::clone(pipeline);
221 self.scope.as_ref().unwrap().spawn(handle_receiver_requests::<
223 fdiagnostics::ArchiveAccessorMarker,
224 >(
225 get_receiver_stream(
226 pipeline.protocol_name(),
227 accessors_dict_id,
228 id_gen,
229 capability_store,
230 )
231 .await,
232 Arc::clone(&accessor_server),
233 Arc::clone(&accessor_pipeline),
234 ));
235 self.scope.as_ref().unwrap().spawn(handle_receiver_requests::<
237 fdiagnostics_host::ArchiveAccessorMarker,
238 >(
239 get_receiver_stream(
240 pipeline.host_protocol_name(),
241 accessors_dict_id,
242 id_gen,
243 capability_store,
244 )
245 .await,
246 Arc::clone(&accessor_server),
247 Arc::clone(&accessor_pipeline),
248 ));
249 }
250
251 PipelinesDictionaryId(accessors_dict_id)
252 }
253}
254
255async fn get_receiver_stream(
256 protocol_name: Cow<'_, str>,
257 accessors_dict_id: u64,
258 id_gen: &sandbox::CapabilityIdGenerator,
259 capability_store: &mut fsandbox::CapabilityStoreProxy,
260) -> fsandbox::ReceiverRequestStream {
261 let (accessor_receiver_client, receiver_stream) =
262 fidl::endpoints::create_request_stream::<fsandbox::ReceiverMarker>();
263 let connector_id = id_gen.next();
264 capability_store
265 .connector_create(connector_id, accessor_receiver_client)
266 .await
267 .unwrap()
268 .unwrap();
269 debug!("Added {protocol_name} to the accessors dictionary.");
270 capability_store
271 .dictionary_insert(
272 accessors_dict_id,
273 &fsandbox::DictionaryItem { key: protocol_name.into_owned(), value: connector_id },
274 )
275 .await
276 .unwrap()
277 .unwrap();
278 receiver_stream
279}
280
281async fn handle_receiver_requests<P>(
282 mut receiver_stream: fsandbox::ReceiverRequestStream,
283 accessor_server: Arc<ArchiveAccessorServer>,
284 pipeline: Arc<Pipeline>,
285) where
286 P: ProtocolMarker,
287 P::RequestStream: ArchiveAccessorTranslator + Send + 'static,
288 <P::RequestStream as ArchiveAccessorTranslator>::InnerDataRequestChannel:
289 ArchiveAccessorWriter + Send,
290{
291 while let Some(request) = receiver_stream.try_next().await.unwrap() {
292 match request {
293 fsandbox::ReceiverRequest::Receive { channel, control_handle: _ } => {
294 debug!("Handling receive request for: {} -> {}", pipeline.name, P::DEBUG_NAME);
295 let server_end = ServerEnd::<P>::new(channel);
296 accessor_server.spawn_server::<P::RequestStream>(
297 Arc::clone(&pipeline),
298 server_end.into_stream(),
299 );
300 }
301 fsandbox::ReceiverRequest::_UnknownMethod { method_type, ordinal, .. } => {
302 warn!(method_type:?, ordinal; "Got unknown interaction on Receiver");
303 }
304 }
305 }
306}