archivist_lib/logs/servers/
log_settings.rs1use crate::logs::error::LogsError;
6use crate::logs::repository::{LogsRepository, STATIC_CONNECTION_ID};
7use fidl::endpoints::DiscoverableProtocolMarker;
8use fidl_fuchsia_diagnostics as fdiagnostics;
9use fuchsia_async as fasync;
10use futures::StreamExt;
11use log::warn;
12use std::sync::Arc;
13
14pub struct LogSettingsServer {
15 logs_repo: Arc<LogsRepository>,
17
18 scope: fasync::Scope,
20}
21
22impl LogSettingsServer {
23 pub fn new(logs_repo: Arc<LogsRepository>, scope: fasync::Scope) -> Self {
24 Self { logs_repo, scope }
25 }
26
27 pub fn spawn(&self, stream: fdiagnostics::LogSettingsRequestStream) {
29 let logs_repo = Arc::clone(&self.logs_repo);
30 self.scope.spawn(async move {
31 if let Err(e) = Self::handle_requests(logs_repo, stream).await {
32 warn!("error handling Log requests: {}", e);
33 }
34 });
35 }
36
37 pub async fn handle_requests(
38 logs_repo: Arc<LogsRepository>,
39 mut stream: fdiagnostics::LogSettingsRequestStream,
40 ) -> Result<(), LogsError> {
41 let connection_id = logs_repo.new_interest_connection();
42 while let Some(request) = stream.next().await {
43 let request = request.map_err(|source| LogsError::HandlingRequests {
44 protocol: fdiagnostics::LogSettingsMarker::PROTOCOL_NAME,
45 source,
46 })?;
47 match request {
48 fidl_fuchsia_diagnostics::LogSettingsRequest::SetComponentInterest {
49 payload,
50 responder,
51 } => {
52 if let Some(selectors) = payload.selectors {
53 let connection_id = if payload.persist.unwrap_or(false) {
54 STATIC_CONNECTION_ID
55 } else {
56 connection_id
57 };
58 logs_repo.update_logs_interest(connection_id, selectors);
59 }
60 responder.send().ok();
61 }
62 }
63 }
64 logs_repo.finish_interest_connection(connection_id);
65
66 Ok(())
67 }
68}
69
70#[cfg(test)]
71mod tests {
72 use super::*;
73 use crate::logs::shared_buffer::create_ring_buffer;
74 use fidl::endpoints::{Proxy, create_proxy_and_stream};
75 use fidl_fuchsia_diagnostics_types::{Interest, Severity};
76 use selectors::{VerboseError, parse_component_selector};
77
78 fn init_repo() -> Arc<LogsRepository> {
79 LogsRepository::new(
80 create_ring_buffer(65536),
81 std::iter::empty(),
82 &Default::default(),
83 fasync::Scope::new(),
84 )
85 }
86
87 #[fuchsia::test]
88 async fn set_component_interest_with_selectors() {
89 let repo = init_repo();
90 let scope = fasync::Scope::new();
91 let server = LogSettingsServer::new(Arc::clone(&repo), scope);
92
93 let (proxy, stream) = create_proxy_and_stream::<fdiagnostics::LogSettingsMarker>();
94 server.spawn(stream);
95
96 let component_selector = parse_component_selector::<VerboseError>("a/b/c").unwrap();
97 let interests = vec![fdiagnostics::LogInterestSelector {
98 selector: component_selector,
99 interest: Interest { min_severity: Some(Severity::Info), ..Default::default() },
100 }];
101
102 proxy
104 .set_component_interest(&fdiagnostics::LogSettingsSetComponentInterestRequest {
105 selectors: Some(interests.clone()),
106 persist: Some(false),
107 ..Default::default()
108 })
109 .await
110 .unwrap();
111
112 proxy
114 .set_component_interest(&fdiagnostics::LogSettingsSetComponentInterestRequest {
115 selectors: Some(interests),
116 persist: Some(true),
117 ..Default::default()
118 })
119 .await
120 .unwrap();
121 }
122
123 #[fuchsia::test]
124 async fn set_component_interest_without_selectors() {
125 let repo = init_repo();
126 let scope = fasync::Scope::new();
127 let server = LogSettingsServer::new(Arc::clone(&repo), scope);
128
129 let (proxy, stream) = create_proxy_and_stream::<fdiagnostics::LogSettingsMarker>();
130 server.spawn(stream);
131
132 proxy
134 .set_component_interest(&fdiagnostics::LogSettingsSetComponentInterestRequest {
135 selectors: None,
136 persist: None,
137 ..Default::default()
138 })
139 .await
140 .unwrap();
141 }
142
143 #[fuchsia::test]
144 async fn stream_error_handling() {
145 let repo = init_repo();
146 let (proxy, stream) = create_proxy_and_stream::<fdiagnostics::LogSettingsMarker>();
147
148 proxy.as_channel().write(&[0xff; 16], &mut []).expect("write invalid bytes");
150
151 let result = LogSettingsServer::handle_requests(repo, stream).await;
152 assert_matches::assert_matches!(result, Err(LogsError::HandlingRequests { .. }));
153 }
154}