1use crate::config::SamplerConfig;
6use crate::project::Project;
7use anyhow::Error as AnyhowError;
8use argh::FromArgs;
9use diagnostics_reader::drain_batch_iterator;
10use fidl::endpoints::{ControlHandle, RequestStream, create_endpoints};
11use fidl_fuchsia_diagnostics as fdiagnostics;
12use fidl_fuchsia_hardware_power_statecontrol::{
13 ShutdownWatcherMarker, ShutdownWatcherRegisterMarker, ShutdownWatcherRequest,
14};
15use fidl_fuchsia_metrics::MetricEventLoggerFactoryMarker;
16use fuchsia_component_client::connect_to_protocol;
17use fuchsia_inspect::component;
18use fuchsia_inspect::health::Reporter;
19use futures::future::{Either, select};
20use futures::stream::{self, StreamExt};
21use inspect_runtime::publish;
22use itertools::Itertools;
23use log::{info, warn};
24use sampler_component_config::Config;
25use std::sync::Arc;
26
27mod config;
28mod error;
29mod project;
30
31#[derive(Debug, Default, FromArgs, PartialEq)]
33#[argh(subcommand, name = "sampler")]
34pub struct Args {}
35
36pub const PROGRAM_NAME: &str = "sampler";
37
38pub async fn main() -> Result<(), AnyhowError> {
39 info!("Sampler starting up");
40 component::health().set_starting_up();
41
42 let _inspect = publish(component::inspector(), Default::default());
43
44 let execution_stats = component::inspector().root().create_child("sampler_executor_stats");
45 let config = SamplerConfig::new(Config::take_from_startup_handle(), &execution_stats)?;
46
47 let sampler = connect_to_protocol::<fdiagnostics::SampleMarker>()?;
48
49 for chunk in &config
50 .sample_data()
51 .into_iter()
52 .chunks(fdiagnostics::MAX_SAMPLE_PARAMETERS_PER_SET as usize)
53 {
54 sampler.set(&fdiagnostics::SampleParameters {
55 data: Some(chunk.collect()),
56 ..Default::default()
57 })?;
58 }
59
60 let (sample_sink_client, sample_sink_server) =
61 create_endpoints::<fdiagnostics::SampleSinkMarker>();
62
63 if let Err(e) = sampler.commit(sample_sink_client).await? {
64 match e {
65 fdiagnostics::ConfigurationError::SamplePeriodTooSmall => {
66 return Err(anyhow::anyhow!(
67 "Configured sample period was too small, indicating a config bug. Exiting."
68 ));
69 }
70 err => warn!(err:?; "Sampler encountered non-fatal error. Review Archivist's logs."),
71 }
72 }
73
74 let metric_logger_factory = connect_to_protocol::<MetricEventLoggerFactoryMarker>()?;
75
76 let mut projects = futures::stream::iter(config.project_configs)
77 .filter_map(|project_config| async {
78 let project_id = *project_config.project_id;
79 let stats = config.stats.projects.get(&project_config.project_id);
80 match Project::new(&metric_logger_factory, project_config, stats).await {
81 Ok(project) => Some(project),
82 Err(e) => {
83 warn!(
84 e:?,
85 project_id;
86 "Sampler failed to configure a project",
87 );
88 None
89 }
90 }
91 })
92 .collect::<Vec<_>>()
93 .await;
94
95 let (shutdown_watcher_client, shutdown_watcher_request_stream) =
96 fidl::endpoints::create_request_stream::<ShutdownWatcherMarker>();
97 let shutdown_watcher_register = connect_to_protocol::<ShutdownWatcherRegisterMarker>()?;
98 shutdown_watcher_register.register_watcher(shutdown_watcher_client).await?;
99
100 let sink_stream = sample_sink_server.into_stream();
101 let sample_sink_control = sink_stream.control_handle();
102 let mut sink_stream = sink_stream.fuse();
103 let mut shutdown_stream = Either::Left(shutdown_watcher_request_stream);
104 let mut shutdown = false;
105 let mut shutdown_responder = None;
106
107 component::health().set_ok();
108
109 loop {
110 match select(shutdown_stream.next(), sink_stream.next()).await {
111 Either::Left((shutdown_request, _)) => match shutdown_request {
112 Some(Ok(ShutdownWatcherRequest::OnShutdown { responder, .. })) => {
113 shutdown = true;
114 shutdown_responder = Some(responder);
115 shutdown_stream = Either::Right(stream::pending());
116 if let Err(e) = sample_sink_control.send_on_now_or_never() {
117 warn!(e:?; "Failed to send on_now_or_never to sample sink");
118 break;
119 }
120 }
121 Some(Ok(ShutdownWatcherRequest::_UnknownMethod { .. })) => {
122 warn!("Sampler encountered unknown method on ShutdownWatcher");
123 }
124 Some(Err(err)) => {
125 warn!(err:?; "Sampler encountered error on ShutdownWatcher, data may be missing");
126 if shutdown {
127 break;
128 }
129 }
130 None => {
131 if shutdown {
132 break;
133 }
134 shutdown_stream = Either::Right(stream::pending());
135 continue;
136 }
137 },
138 Either::Right((event, _)) => {
139 let event = match event {
140 Some(Ok(event)) => event,
141 Some(Err(err)) => {
142 warn!(err:?; "Sample sink stream encountered an error");
143 break;
144 }
145 None => break,
146 };
147
148 handle_sample_sink_request(event, shutdown, &mut projects).await;
149
150 if shutdown {
151 break;
152 }
153 }
154 }
155 }
156
157 if let Some(responder) = shutdown_responder
158 && let Err(err) = responder.send()
159 {
160 warn!(err:?; "Failed to send ShutdownWatcher response; shutdown coordinator may have timed out");
161 }
162
163 Ok(())
164}
165
166async fn handle_sample_sink_request(
167 event: fdiagnostics::SampleSinkRequest,
168 shutdown: bool,
169 projects: &mut [Project<'_>],
170) {
171 match event {
172 fdiagnostics::SampleSinkRequest::OnSampleReadied {
173 event:
174 fdiagnostics::SampleSinkResult::Ready(fdiagnostics::SampleReady {
175 batch_iter: Some(batch_iter),
176 seconds_since_start: Some(seconds_since_start),
177 ..
178 }),
179 control_handle: _control_handle,
180 } => {
181 let data = drain_batch_iterator::<diagnostics_data::InspectData>(Arc::new(
182 batch_iter.into_proxy(),
183 ))
184 .filter_map(|v| async {
185 match v {
186 Ok(v) => Some(v),
187 Err(e) => {
188 warn!(e:?; "Failed to read some Inspect data; skipping");
189 None
190 }
191 }
192 })
193 .collect::<Vec<_>>()
194 .await;
195
196 let seconds_since_start = if shutdown {
197 None
198 } else {
199 Some(zx::MonotonicDuration::from_seconds(seconds_since_start))
200 };
201
202 for project in projects {
203 if let Err(e) = project.log(&data, seconds_since_start).await {
204 warn!(e:?; "Project failed to log");
205 }
206
207 }
211 }
212 fdiagnostics::SampleSinkRequest::OnSampleReadied {
213 event:
214 fdiagnostics::SampleSinkResult::Ready(fdiagnostics::SampleReady {
215 batch_iter,
216 seconds_since_start,
217 ..
218 }),
219 control_handle,
220 } => {
221 control_handle.shutdown();
222 warn!(
223 batch_iter:?, seconds_since_start:?;
224 "Sample server sent Ready but crucial fields were None"
225 );
226 }
227 fdiagnostics::SampleSinkRequest::OnSampleReadied {
228 event: fdiagnostics::SampleSinkResult::Error(e),
229 ..
230 } => {
231 warn!(e:?; "Sample server sent an error, data may be missing");
232 }
233 fdiagnostics::SampleSinkRequest::OnSampleReadied {
234 event: fdiagnostics::SampleSinkResult::__SourceBreaking { .. },
235 control_handle,
236 }
237 | fdiagnostics::SampleSinkRequest::_UnknownMethod { control_handle, .. } => {
238 control_handle.shutdown();
239 warn!("Sample server sent a source-breaking or unknown event")
240 }
241 }
242}