1use anyhow::Context as _;
11use fidl_fuchsia_net as fnet;
12use fidl_fuchsia_net_policy_socketproxy as fnp_socketproxy;
13use fidl_fuchsia_posix_socket::{self as fposix_socket, OptionalUint32};
14use fuchsia_async as fasync;
15use fuchsia_component::server::{ServiceFs, ServiceFsDir};
16use fuchsia_inspect::health::Reporter;
17use fuchsia_inspect_derive::{Inspect, WithInspect as _};
18use futures::StreamExt as _;
19use futures::channel::mpsc;
20use futures::lock::Mutex;
21use log::error;
22use std::sync::Arc;
23
24mod mark_watcher;
25pub mod registry;
26mod socket_provider;
27
28pub use registry::NetworkRegistryError;
29
30#[derive(Copy, Clone, Debug)]
31struct SocketMarks {
32 mark_1: OptionalUint32,
33 mark_2: OptionalUint32,
34}
35
36impl From<SocketMarks> for fnet::Marks {
37 fn from(SocketMarks { mark_1, mark_2 }: SocketMarks) -> Self {
38 let into_option_u32 = |opt| match opt {
39 OptionalUint32::Unset(fposix_socket::Empty) => None,
40 OptionalUint32::Value(val) => Some(val),
41 };
42 Self {
43 mark_1: into_option_u32(mark_1),
44 mark_2: into_option_u32(mark_2),
45 __source_breaking: fidl::marker::SourceBreaking,
46 }
47 }
48}
49
50impl SocketMarks {
51 fn has_value(&self) -> bool {
52 match (self.mark_1, self.mark_2) {
53 (OptionalUint32::Value(_), _) => true,
54 (_, OptionalUint32::Value(_)) => true,
55 _ => false,
56 }
57 }
58
59 fn set_mark(&mut self, domain: fnet::MarkDomain, value: Option<u32>) {
60 let value = match value {
61 Some(value) => fposix_socket::OptionalUint32::Value(value),
62 None => fposix_socket::OptionalUint32::Unset(fposix_socket::Empty),
63 };
64
65 match domain {
66 fnet::MarkDomain::Mark1 => self.mark_1 = value,
67 fnet::MarkDomain::Mark2 => self.mark_2 = value,
68 }
69 }
70}
71
72impl Default for SocketMarks {
73 fn default() -> Self {
74 Self {
75 mark_1: OptionalUint32::Unset(fposix_socket::Empty),
76 mark_2: OptionalUint32::Unset(fposix_socket::Empty),
77 }
78 }
79}
80
81#[derive(Inspect)]
82struct SocketProxy {
83 registry: registry::Registry,
84 socket_provider: socket_provider::SocketProvider,
85}
86
87impl SocketProxy {
88 fn new(
89 forwarder_tx: mpsc::Sender<crate::registry::NetworkRegistryRequest>,
90 ) -> Result<Self, anyhow::Error> {
91 let mark = Arc::new(Mutex::new(SocketMarks::default()));
92 Ok(Self {
93 registry: registry::Registry::new(mark.clone(), forwarder_tx)
94 .context("while creating registry")?,
95 socket_provider: socket_provider::SocketProvider::new(mark),
96 })
97 }
98}
99
100enum IncomingService {
101 StarnixNetworks(fnp_socketproxy::StarnixNetworksRequestStream),
102 PosixSocket(fidl_fuchsia_posix_socket::ProviderRequestStream),
103 PosixSocketRaw(fidl_fuchsia_posix_socket_raw::ProviderRequestStream),
104}
105
106pub async fn run() -> Result<(), anyhow::Error> {
108 fuchsia_inspect::component::health().set_starting_up();
109
110 let inspector = fuchsia_inspect::component::inspector();
111 let _inspect_server_task =
112 inspect_runtime::publish(inspector, inspect_runtime::PublishOptions::default());
113
114 let (forwarder_tx, forwarder_rx) = mpsc::channel(50);
118 let mut request_forwarder = registry::RequestForwarder::new(forwarder_rx)?;
119
120 let proxy = Arc::new(SocketProxy::new(forwarder_tx)?.with_inspect(inspector.root(), "root")?);
121
122 let mut fs = ServiceFs::new_local();
123 let _: &mut ServiceFsDir<'_, _> = fs
124 .dir("svc")
125 .add_fidl_service(IncomingService::StarnixNetworks)
126 .add_fidl_service(IncomingService::PosixSocket)
127 .add_fidl_service(IncomingService::PosixSocketRaw);
128
129 let _: &mut ServiceFs<_> = fs.take_and_serve_directory_handle()?;
130
131 fuchsia_inspect::component::health().set_ok();
132
133 let proxy_for_service = Arc::clone(&proxy);
134 let service_fut = fs.for_each_concurrent(100, move |service| {
135 let proxy = Arc::clone(&proxy_for_service);
136 async move {
137 match service {
138 IncomingService::StarnixNetworks(stream) => {
139 proxy.registry.run_starnix(stream).await
140 }
141 IncomingService::PosixSocket(stream) => proxy.socket_provider.run(stream).await,
142 IncomingService::PosixSocketRaw(stream) => {
143 proxy.socket_provider.run_raw(stream).await
144 }
145 }
146 .unwrap_or_else(|e| error!("{e:?}"))
147 }
148 });
149
150 let scope = fasync::Scope::new();
151
152 let _ = scope.spawn_local(async move {
153 let res = request_forwarder.run().await;
154 error!("RequestForwarder future has terminated: {res:?}");
155 });
156
157 let proxy_clone = Arc::clone(&proxy);
158 let _ = scope.spawn_local(async move {
159 mark_watcher::watch_properties(proxy_clone).await;
160 });
161
162 let _ = scope.spawn_local(async move {
163 service_fut.await;
164 error!("The main services future has terminated. It should never terminate");
165 fasync::Scope::current().abort().await;
167 });
168
169 scope.join().await;
170
171 Ok(())
172}