archivist_lib/logs/servers/
log.rs1use 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 logs_repo: Arc<LogsRepository>,
20
21 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 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 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 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}