Skip to main content

archivist_lib/logs/servers/
log.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::listener::Listener;
7use crate::logs::repository::LogsRepository;
8use fidl::endpoints::DiscoverableProtocolMarker;
9use fidl_fuchsia_diagnostics::StreamMode;
10use fidl_fuchsia_logger as flogger;
11use fuchsia_async as fasync;
12use futures::StreamExt;
13use log::warn;
14use std::pin::pin;
15use std::sync::Arc;
16
17pub struct LogServer {
18    /// The repository holding the logs.
19    logs_repo: Arc<LogsRepository>,
20
21    /// Scope in which we spawn all of the server tasks.
22    scope: fasync::Scope,
23}
24
25impl LogServer {
26    pub fn new(logs_repo: Arc<LogsRepository>, scope: fasync::Scope) -> Self {
27        Self { logs_repo, scope }
28    }
29
30    /// Spawn a task to handle requests from components reading the shared log.
31    pub fn spawn(&self, stream: flogger::LogRequestStream) {
32        let logs_repo = Arc::clone(&self.logs_repo);
33        let scope = self.scope.to_handle();
34        self.scope.spawn(async move {
35            if let Err(e) = Self::handle_requests(logs_repo, stream, scope).await {
36                warn!("error handling Log requests: {}", e);
37            }
38        });
39    }
40
41    /// Handle requests to `fuchsia.logger.Log`. All request types read the
42    /// whole backlog from memory, `DumpLogs(Safe)` stops listening after that.
43    async fn handle_requests(
44        logs_repo: Arc<LogsRepository>,
45        mut stream: flogger::LogRequestStream,
46        scope: fasync::ScopeHandle,
47    ) -> Result<(), LogsError> {
48        let connection_id = logs_repo.new_interest_connection();
49        while let Some(request) = stream.next().await {
50            let request = request.map_err(|source| LogsError::HandlingRequests {
51                protocol: flogger::LogMarker::PROTOCOL_NAME,
52                source,
53            })?;
54            let listener = match request {
55                flogger::LogRequest::ListenSafe { log_listener, options, .. } => {
56                    Listener::new(log_listener, options)?
57                }
58            };
59            let logs =
60                logs_repo.logs_cursor(StreamMode::SnapshotThenSubscribe, Vec::new()).map(Arc::new);
61            scope.spawn(async move {
62                listener.run(pin!(logs)).await;
63            });
64        }
65        logs_repo.finish_interest_connection(connection_id);
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
76    fn init_repo() -> Arc<LogsRepository> {
77        LogsRepository::new(
78            create_ring_buffer(65536),
79            std::iter::empty(),
80            &Default::default(),
81            fasync::Scope::new(),
82        )
83    }
84
85    #[fuchsia::test]
86    async fn listen_safe_and_disconnect() {
87        let repo = init_repo();
88        let scope = fasync::Scope::new();
89        let server = LogServer::new(Arc::clone(&repo), scope);
90
91        let (proxy, stream) = create_proxy_and_stream::<flogger::LogMarker>();
92        server.spawn(stream);
93
94        let (client_end, _server_end) =
95            fidl::endpoints::create_endpoints::<flogger::LogListenerSafeMarker>();
96        proxy.listen_safe(client_end, None).unwrap();
97
98        drop(proxy);
99        fasync::Timer::new(std::time::Duration::from_millis(10)).await;
100    }
101
102    #[fuchsia::test]
103    async fn stream_error_handling() {
104        let repo = init_repo();
105        let scope = fasync::Scope::new();
106        let (proxy, stream) = create_proxy_and_stream::<flogger::LogMarker>();
107
108        // Write invalid bytes to cause a stream error.
109        proxy.as_channel().write(&[0xff; 16], &mut []).expect("write invalid bytes");
110
111        let result = LogServer::handle_requests(repo, stream, scope.to_handle()).await;
112        assert_matches::assert_matches!(result, Err(LogsError::HandlingRequests { .. }));
113    }
114}