Skip to main content

archivist_lib/logs/servers/
log_settings.rs

1// Copyright 2022 The Fuchsia Authors. All rights reserved.
2// Use of this source code is governed by a BSD-style license that can be
3// found in the LICENSE file.
4
5use 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    /// The repository holding the logs.
16    logs_repo: Arc<LogsRepository>,
17
18    /// Scope holding all of the server Tasks.
19    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    /// Spawn a task to handle requests from components reading the shared log.
28    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        // Test non-persistent
103        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        // Test persistent
113        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        // Test with selectors: None
133        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        // Write invalid header to trigger decode failure in stream.
149        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}