Skip to main content

dhcpv4_server/
main.rs

1// Copyright 2018 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 anyhow::{Context as _, Error};
6use dhcpv4::configuration;
7use dhcpv4::protocol::{CLIENT_PORT, Message, SERVER_PORT};
8use dhcpv4::server::{
9    DEFAULT_STASH_ID, DataStore, ResponseTarget, Server, ServerAction, ServerDispatcher,
10    ServerError,
11};
12use dhcpv4::stash::Stash;
13use fuchsia_async::net::UdpSocket;
14use fuchsia_async::{self as fasync};
15use fuchsia_component::server::{ServiceFs, ServiceFsDir};
16use futures::{Future, SinkExt as _, StreamExt as _, TryFutureExt as _, TryStreamExt as _};
17use log::{debug, error, info, warn};
18use net_declare::net::prefix_length_v4;
19use net_types::ethernet::Mac;
20use packet::serialize::InnerPacketBuilder;
21use packet::{NestableSerializer as _, NoOpSerializationContext, Serializer};
22use packet_formats::ipv4::Ipv4PacketBuilder;
23use packet_formats::udp::UdpPacketBuilder;
24use sockaddr::IntoSockAddr as _;
25use std::cell::RefCell;
26use std::collections::HashMap;
27use std::net::{IpAddr, Ipv4Addr, SocketAddr};
28
29/// A buffer size in excess of the maximum allowable DHCP message size.
30const BUF_SZ: usize = 1024;
31
32enum IncomingService {
33    Server(fidl_fuchsia_net_dhcp::Server_RequestStream),
34}
35
36const DEFAULT_LEASE_DURATION_SECONDS: u32 = 24 * 60 * 60;
37
38fn default_parameters() -> configuration::ServerParameters {
39    configuration::ServerParameters {
40        server_ips: vec![],
41        lease_length: dhcpv4::configuration::LeaseLength {
42            default_seconds: DEFAULT_LEASE_DURATION_SECONDS,
43            max_seconds: DEFAULT_LEASE_DURATION_SECONDS,
44        },
45        managed_addrs: dhcpv4::configuration::ManagedAddresses {
46            mask: configuration::SubnetMask::new(prefix_length_v4!(0)),
47            pool_range_start: Ipv4Addr::UNSPECIFIED,
48            pool_range_stop: Ipv4Addr::UNSPECIFIED,
49        },
50        permitted_macs: dhcpv4::configuration::PermittedMacs(vec![]),
51        static_assignments: dhcpv4::configuration::StaticAssignments(
52            std::collections::hash_map::HashMap::new(),
53        ),
54        arp_probe: false,
55        bound_device_names: vec![],
56    }
57}
58
59/// dhcpd is the Fuchsia DHCPv4 server.
60#[derive(argh::FromArgs)]
61struct Args {
62    /// enables storage of dhcpd lease and configuration state to persistent storage
63    #[argh(switch)]
64    persistent: bool,
65}
66
67#[fuchsia::main()]
68pub async fn main() -> Result<(), Error> {
69    info!("starting");
70
71    let Args { persistent } = argh::from_env();
72    info!("persistent={}", persistent);
73    if persistent {
74        let stash = Stash::new(DEFAULT_STASH_ID).context("failed to instantiate stash")?;
75        // The server parameters and the client records must be consistent with one another in
76        // order to ensure correct server operation. The records cannot be consistent with default
77        // parameters, so if parameters fail to load from the stash, then the records should
78        // default to empty.
79        let (params, options, records) = match stash.load_parameters().await {
80            Ok(params) => {
81                let options = stash.load_options().await.unwrap_or_else(|e| {
82                    warn!("failed to load options from stash: {:?}", e);
83                    HashMap::new()
84                });
85                let records = stash.load_client_records().await.unwrap_or_else(|e| {
86                    warn!("failed to load client records from stash: {:?}", e);
87                    HashMap::new()
88                });
89                (params, options, records)
90            }
91            Err(e) => {
92                warn!("failed to load parameters from stash: {:?}", e);
93                (default_parameters(), HashMap::new(), HashMap::new())
94            }
95        };
96        let server = match Server::new_from_state(stash.clone(), params, options, records) {
97            Ok(v) => v,
98            Err(e) => {
99                warn!("failed to create server from persistent state: {}", e);
100                Server::new(Some(stash), default_parameters())
101            }
102        };
103        Ok(run(server).await?)
104    } else {
105        Ok(run(Server::<Stash>::new(None, default_parameters())).await?)
106    }
107}
108
109async fn run<DS: DataStore>(server: Server<DS>) -> Result<(), Error> {
110    let server = RefCell::new(ServerDispatcherRuntime::new(server));
111
112    let mut fs = ServiceFs::new_local();
113    let _: &mut ServiceFsDir<'_, _> = fs.dir("svc").add_fidl_service(IncomingService::Server);
114    let _: &mut ServiceFs<_> = fs
115        .take_and_serve_directory_handle()
116        .context("service fs failed to take and serve directory handle")?;
117
118    let (mut socket_sink, socket_stream) =
119        futures::channel::mpsc::channel::<ServerSocketCollection<UdpSocket>>(1);
120
121    // Attempt to enable the server on startup.
122    // NOTE(brunodalbo): Enabling the server on startup should be an explicit
123    // configuration loaded from default configs and stash. For now, just mimic
124    // existing behavior and try to enable. It'll fail if we don't have a valid
125    // configuration from stash/config.
126    match server.borrow_mut().enable() {
127        Ok(None) => unreachable!("server can't be enabled already"),
128        Ok(Some(socket_collection)) => {
129            // Sending here should never fail; we just created the stream above.
130            socket_sink.try_send(socket_collection)?;
131        }
132        Err(e @ zx::Status::INVALID_ARGS) => {
133            info!("server not configured for serving leases: {:?}", e)
134        }
135        Err(e) => warn!("could not enable server on startup: {:?}", e),
136    }
137
138    let admin_fut =
139        fs.then(futures::future::ok).try_for_each_concurrent(None, |incoming_service| async {
140            match incoming_service {
141                IncomingService::Server(stream) => {
142                    run_server(stream, &server, &default_parameters(), socket_sink.clone())
143                        .inspect_err(|e| warn!("run_server failed: {:?}", e))
144                        .await?;
145                    Ok(())
146                }
147            }
148        });
149
150    let server_fut = define_running_server_fut(&server, socket_stream);
151
152    info!("running");
153    let ((), ()) = futures::try_join!(server_fut, admin_fut)?;
154
155    Ok(())
156}
157
158trait SocketServerDispatcher: ServerDispatcher {
159    type Socket;
160
161    fn create_socket(name: &str, src: Ipv4Addr) -> std::io::Result<Self::Socket>;
162    fn dispatch_message(&mut self, msg: Message) -> Result<ServerAction, ServerError>;
163    fn create_sockets(
164        params: &configuration::ServerParameters,
165    ) -> std::io::Result<Vec<SocketWithId<Self::Socket>>>;
166}
167
168impl<DS: DataStore> SocketServerDispatcher for Server<DS> {
169    type Socket = UdpSocket;
170
171    fn create_socket(name: &str, src: Ipv4Addr) -> std::io::Result<Self::Socket> {
172        let socket = socket2::Socket::new(
173            socket2::Domain::IPV4,
174            socket2::Type::DGRAM,
175            Some(socket2::Protocol::UDP),
176        )?;
177        // Since dhcpd may listen to multiple interfaces, we must enable
178        // SO_REUSEPORT so that binding the same (address, port) pair to each
179        // interface can still succeed.
180        socket.set_reuse_port(true)?;
181        socket.bind_device(Some(name.as_bytes()))?;
182        info!("socket bound to device {}", name);
183        socket.set_broadcast(true)?;
184        socket.bind(&SocketAddr::new(IpAddr::V4(src), SERVER_PORT.into()).into())?;
185        Ok(UdpSocket::from_socket(socket.into())?)
186    }
187
188    fn dispatch_message(&mut self, msg: Message) -> Result<ServerAction, ServerError> {
189        self.dispatch(msg)
190    }
191
192    fn create_sockets(
193        params: &configuration::ServerParameters,
194    ) -> std::io::Result<Vec<SocketWithId<Self::Socket>>> {
195        let configuration::ServerParameters { bound_device_names, .. } = params;
196        bound_device_names
197            .iter()
198            .map(|name| {
199                let iface_id =
200                    fuchsia_nix::net::if_::if_nametoindex(name.as_str()).map_err(|e| {
201                        let e: std::io::Error = e.into();
202                        e
203                    })?;
204                let socket = Self::create_socket(name, Ipv4Addr::UNSPECIFIED)?;
205                Ok(SocketWithId { socket, iface_id: iface_id.into() })
206            })
207            .collect()
208    }
209}
210
211/// A wrapper around a [`ServerDispatcher`] that keeps information about the
212/// server status through a [`futures::future::AbortHandle`].
213struct ServerDispatcherRuntime<S> {
214    abort_handle: Option<futures::future::AbortHandle>,
215    server: S,
216}
217
218impl<S> std::ops::Deref for ServerDispatcherRuntime<S> {
219    type Target = S;
220
221    fn deref(&self) -> &Self::Target {
222        &self.server
223    }
224}
225
226impl<S> std::ops::DerefMut for ServerDispatcherRuntime<S> {
227    fn deref_mut(&mut self) -> &mut Self::Target {
228        &mut self.server
229    }
230}
231
232impl<S: SocketServerDispatcher> ServerDispatcherRuntime<S> {
233    /// Creates a new runtime with `server`.
234    fn new(server: S) -> Self {
235        Self { abort_handle: None, server }
236    }
237
238    /// Disables the server.
239    ///
240    /// `disable` will cancel the previous
241    /// [`futures::future::AbortRegistration`] returned by `enable`.
242    ///
243    /// If the server is already disabled, `disable` is a no-op.
244    fn disable(&mut self) {
245        if let Some(abort_handle) = self.abort_handle.take() {
246            abort_handle.abort();
247        }
248    }
249
250    /// Enables the server.
251    ///
252    /// Attempts to enable the server, returning a new
253    /// [`ServerSocketCollection`] on success. The returned collection contains
254    /// the list of sockets where the server can listen on and an abort
255    /// registration that is used to cancel the future that listen on the
256    /// sockets when [`ServerDispatcherRuntime::disable`] is called.
257    ///
258    /// Returns an error if the server couldn't be started or if the closure
259    /// fails, maintaining the server in the disabled state.
260    ///
261    /// If the server is already enabled, `enable` returns `Ok(None)`.
262    fn enable(&mut self) -> Result<Option<ServerSocketCollection<S::Socket>>, zx::Status> {
263        if self.abort_handle.is_some() {
264            // Server already running.
265            return Ok(None);
266        }
267        let params = self.server.try_validate_parameters()?;
268        // Provide the closure with an AbortRegistration and a ref to
269        // parameters.
270        let (abort_handle, abort_registration) = futures::future::AbortHandle::new_pair();
271
272        let sockets = S::create_sockets(params).map_err(|e| {
273            match e.raw_os_error() {
274                // A short-lived SoftAP interface may be, and frequently is, torn down prior to the
275                // full instantiation of its associated dhcpd component. Consequently, binding to
276                // the SoftAP interface name will fail with ENODEV. However, such a failure is
277                // normal and expected under those circumstances.
278                Some(libc::ENODEV) => {
279                    warn!("Failed to create server sockets: {}", e)
280                }
281                Some(_) | None => error!("Failed to create server sockets: {}", e),
282            };
283            zx::Status::IO
284        })?;
285        if sockets.is_empty() {
286            error!("No sockets to run server on");
287            return Err(zx::Status::INVALID_ARGS);
288        }
289        self.abort_handle = Some(abort_handle);
290        Ok(Some(ServerSocketCollection { sockets, abort_registration }))
291    }
292
293    /// Returns `true` if the server is enabled.
294    fn enabled(&self) -> bool {
295        self.abort_handle.is_some()
296    }
297
298    /// Runs the closure `f` only if the server is currently disabled.
299    ///
300    /// Returns `BAD_STATE` error otherwise.
301    fn if_disabled<R, F: FnOnce(&mut S) -> Result<R, zx::Status>>(
302        &mut self,
303        f: F,
304    ) -> Result<R, zx::Status> {
305        if self.abort_handle.is_none() { f(&mut self.server) } else { Err(zx::Status::BAD_STATE) }
306    }
307}
308
309#[derive(Debug, PartialEq)]
310struct SocketWithId<S> {
311    socket: S,
312    iface_id: u64,
313}
314
315/// Helper struct to handle buffer data from sockets.
316struct MessageHandler<'a, S: SocketServerDispatcher> {
317    server: &'a RefCell<ServerDispatcherRuntime<S>>,
318}
319
320impl<'a, S: SocketServerDispatcher> MessageHandler<'a, S> {
321    /// Creates a new `MessageHandler` for `server`.
322    fn new(server: &'a RefCell<ServerDispatcherRuntime<S>>) -> Self {
323        Self { server }
324    }
325
326    /// Handles `buf` from `sender`.
327    ///
328    /// Returns `Ok(Some(sock, msg, dst))` if `msg` must be sent to `dst`
329    /// over `sock`.
330    ///
331    /// Returns `Ok(None)` if no action is required and the handler is ready to
332    /// receive more messages.
333    ///
334    /// Returns `Err` if an unrecoverable error occurs and the server must stop
335    /// serving.
336    fn handle_from_sender(
337        &mut self,
338        buf: &[u8],
339        mut sender: std::net::SocketAddrV4,
340    ) -> Result<Option<(std::net::SocketAddrV4, Message, Option<Mac>)>, Error> {
341        let msg = match Message::from_buffer(buf) {
342            Ok(msg) => {
343                debug!("parsed message from {}: {:?}", sender, msg);
344                msg
345            }
346            Err(e) => {
347                warn!("failed to parse message from {}: {}", sender, e);
348                return Ok(None);
349            }
350        };
351
352        let typ = msg.get_dhcp_type();
353        if sender.ip().is_unspecified() {
354            info!("processing {:?} from {}", typ, msg.chaddr);
355        } else {
356            info!("processing {:?} from {}", typ, sender);
357        }
358
359        // This call should not block because the server is single-threaded.
360        let result = self.server.borrow_mut().dispatch_message(msg);
361        match result {
362            Err(e) => {
363                warn!("error processing client message: {:?}", e);
364                Ok(None)
365            }
366            Ok(ServerAction::AddressRelease(addr)) => {
367                info!("released address: {}", addr);
368                Ok(None)
369            }
370            Ok(ServerAction::AddressDecline(addr)) => {
371                info!("allocated address: {}", addr);
372                Ok(None)
373            }
374            Ok(ServerAction::SendResponse(message, dst)) => {
375                debug!("generated response: {:?}", message);
376
377                let typ = message.get_dhcp_type();
378                // Check if server returned an explicit destination ip.
379                let (addr, chaddr) = match dst {
380                    ResponseTarget::Broadcast => {
381                        info!("sending {:?} to {}", typ, Ipv4Addr::BROADCAST);
382                        (Ipv4Addr::BROADCAST, None)
383                    }
384                    ResponseTarget::Unicast(addr, None) => {
385                        info!("sending {:?} to {}", typ, addr);
386                        (addr, None)
387                    }
388                    ResponseTarget::Unicast(addr, Some(chaddr)) => {
389                        info!("sending {:?} to ip {} chaddr {}", typ, addr, chaddr);
390                        (addr, Some(chaddr))
391                    }
392                };
393                sender.set_ip(addr);
394                Ok(Some((sender, message, chaddr)))
395            }
396        }
397    }
398}
399
400async fn define_msg_handling_loop_future<DS: DataStore>(
401    sock: SocketWithId<<Server<DS> as SocketServerDispatcher>::Socket>,
402    server: &RefCell<ServerDispatcherRuntime<Server<DS>>>,
403) -> Result<!, Error> {
404    let SocketWithId { socket, iface_id } = sock;
405    let mut handler = MessageHandler::new(server);
406    let mut buf = vec![0u8; BUF_SZ];
407    loop {
408        let (received, sender) =
409            socket.recv_from(&mut buf).await.context("failed to read from socket")?;
410        let sender = match sender {
411            std::net::SocketAddr::V4(sender) => sender,
412            std::net::SocketAddr::V6(sender) => {
413                return Err(anyhow::anyhow!(
414                    "IPv4 socket received datagram from IPv6 sender: {}",
415                    sender
416                ));
417            }
418        };
419        if let Some((dst, msg, chaddr)) = handler
420            .handle_from_sender(&buf[..received], sender)
421            .context("failed to handle buffer")?
422        {
423            let chaddr = if let Some(chaddr) = chaddr {
424                chaddr
425            } else {
426                let response = msg.serialize();
427                let sent = socket
428                    .send_to(&response, SocketAddr::V4(dst))
429                    .await
430                    .context("unable to send response")?;
431                if sent != response.len() {
432                    return Err(anyhow::anyhow!(
433                        "sent {} bytes for a message of size {}",
434                        sent,
435                        response.len()
436                    ));
437                }
438                info!("response sent to {}: {} bytes", dst, sent);
439                continue;
440            };
441            // Packet sockets are necessary here because the on-device Netstack does
442            // not yet have a relation linking the `chaddr` MAC address to an IP address.
443            let dst_ip: net_types::ip::Ipv4Addr = (*dst.ip()).into();
444            // Prefer the ServerIdentifier if set; otherwise use the socket's local address.
445            let src_ip: net_types::ip::Ipv4Addr = msg
446                .options
447                .iter()
448                .find_map(|opt| match opt {
449                    dhcpv4::protocol::DhcpOption::ServerIdentifier(addr) => Some(addr.clone()),
450                    _ => None,
451                })
452                // TODO(https://fxbug.dev/42056628): Eliminate this panic.
453                .expect("expect server identifier is always present")
454                .into();
455            let response = msg.serialize();
456            let udp_builder = UdpPacketBuilder::new(src_ip, dst_ip, Some(SERVER_PORT), CLIENT_PORT);
457            // Use the default TTL shared across UNIX systems.
458            const TTL: u8 = 64;
459            let ipv4_builder = Ipv4PacketBuilder::new(
460                src_ip,
461                dst_ip,
462                TTL,
463                packet_formats::ip::Ipv4Proto::Proto(packet_formats::ip::IpProto::Udp),
464            );
465            let packet = response
466                .into_serializer()
467                .wrap_in(udp_builder)
468                .wrap_in(ipv4_builder)
469                .serialize_vec_outer(&mut NoOpSerializationContext)
470                .expect("serialize packet failed")
471                .unwrap_b();
472
473            let mut sll_addr = [0; 8];
474            (&mut sll_addr[..chaddr.bytes().len()]).copy_from_slice(&chaddr.bytes());
475            let sockaddr_ll = libc::sockaddr_ll {
476                sll_family: libc::AF_PACKET.try_into().expect("convert sll_family failed"),
477                sll_ifindex: iface_id.try_into().expect("convert sll_ifindex failed"),
478                // Network order is big endian.
479                sll_protocol: u16::try_from(libc::ETH_P_IP)
480                    .expect("convert ETH_P_IP failed")
481                    .to_be(),
482                sll_halen: chaddr.bytes().len().try_into().expect("convert chaddr size failed"),
483                sll_addr: sll_addr,
484                sll_hatype: 0,
485                sll_pkttype: 0,
486            };
487
488            // Create the packet socket without binding to a protocol so that
489            // the packet socket is not registered for RX. This desirable since
490            // the socket is only used to send packets and receiving packets
491            // on a packet socket is not free.
492            let socket = socket2::Socket::new(
493                socket2::Domain::PACKET,
494                socket2::Type::DGRAM,
495                None, /* protocol */
496            )
497            .context("create packet socket failed")?;
498
499            let socket = fasync::net::DatagramSocket::new_from_socket(socket)
500                .context("failed to wrap into fuchsia-async DatagramSocket")?;
501
502            let sent = socket
503                .send_to(packet.as_ref(), sockaddr_ll.into_sockaddr())
504                .await
505                .context("unable to send response")?;
506            if sent != packet.as_ref().len() {
507                return Err(anyhow::anyhow!(
508                    "sent {} bytes for a packet of size {}",
509                    sent,
510                    packet.as_ref().len()
511                ));
512            }
513            info!("response sent to {}: {} bytes", dst, sent);
514        }
515    }
516}
517
518fn define_running_server_fut<'a, S, DS>(
519    server: &'a RefCell<ServerDispatcherRuntime<Server<DS>>>,
520    socket_stream: S,
521) -> impl Future<Output = Result<(), Error>> + 'a
522where
523    S: futures::Stream<
524            Item = ServerSocketCollection<<Server<Stash> as SocketServerDispatcher>::Socket>,
525        > + 'static,
526    DS: DataStore,
527{
528    socket_stream.map(Ok).try_for_each(move |socket_collection| async move {
529        let ServerSocketCollection { sockets, abort_registration } = socket_collection;
530        let msg_loops = futures::future::try_join_all(
531            sockets.into_iter().map(|sock| define_msg_handling_loop_future(sock, server)),
532        );
533
534        info!("Server starting");
535        match futures::future::Abortable::new(msg_loops, abort_registration).await {
536            Ok(Ok(v)) => {
537                let _: Vec<!> = v;
538                Err(anyhow::anyhow!("Server futures finished unexpectedly"))
539            }
540            Ok(Err(error)) => {
541                // There was an error handling the server sockets. Disable the
542                // server.
543                error!("Server encountered an error: {:?}. Stopping server.", error);
544                server.borrow_mut().disable();
545                Ok(())
546            }
547            Err(futures::future::Aborted {}) => {
548                info!("Server stopped");
549                Ok(())
550            }
551        }
552    })
553}
554
555struct ServerSocketCollection<S> {
556    sockets: Vec<SocketWithId<S>>,
557    abort_registration: futures::future::AbortRegistration,
558}
559
560async fn run_server<S, C>(
561    stream: fidl_fuchsia_net_dhcp::Server_RequestStream,
562    server: &RefCell<ServerDispatcherRuntime<S>>,
563    default_params: &dhcpv4::configuration::ServerParameters,
564    socket_sink: C,
565) -> Result<(), fidl::Error>
566where
567    S: SocketServerDispatcher,
568    C: futures::sink::Sink<ServerSocketCollection<S::Socket>> + Unpin,
569    C::Error: std::fmt::Debug,
570{
571    stream
572        .try_fold(socket_sink, |mut socket_sink, request| async move {
573            match request {
574                fidl_fuchsia_net_dhcp::Server_Request::StartServing { responder } => {
575                    let enable_result = server.borrow_mut().enable();
576
577                    let result = match enable_result {
578                        Ok(Some(socket_collection)) => {
579                            socket_sink.send(socket_collection).await.map_err(|e| {
580                                error!("Failed to send sockets to sink: {:?}", e);
581                                // Disable the server again to keep a consistent state.
582                                server.borrow_mut().disable();
583                                zx::Status::INTERNAL
584                            })
585                        }
586                        Ok(None) => {
587                            info!("Server already running");
588                            Ok(())
589                        }
590                        Err(status) => Err(status),
591                    };
592
593                    responder.send(result.map_err(zx::Status::into_raw))
594                }
595                fidl_fuchsia_net_dhcp::Server_Request::StopServing { responder } => {
596                    server.borrow_mut().disable();
597                    responder.send()
598                }
599                fidl_fuchsia_net_dhcp::Server_Request::IsServing { responder } => {
600                    responder.send(server.borrow().enabled())
601                }
602                fidl_fuchsia_net_dhcp::Server_Request::GetOption { code: c, responder: r } => r
603                    .send(
604                        server.borrow().dispatch_get_option(c).as_ref().map_err(|e| e.into_raw()),
605                    ),
606                fidl_fuchsia_net_dhcp::Server_Request::GetParameter { name: n, responder: r } => {
607                    let response = server.borrow().dispatch_get_parameter(n);
608                    r.send(response.as_ref().map_err(|e| e.into_raw()))
609                }
610                fidl_fuchsia_net_dhcp::Server_Request::SetOption { value: v, responder: r } => {
611                    r.send(server.borrow_mut().dispatch_set_option(v).map_err(|e| e.into_raw()))
612                }
613                fidl_fuchsia_net_dhcp::Server_Request::SetParameter { value: v, responder: r } => r
614                    .send(
615                        server
616                            .borrow_mut()
617                            .if_disabled(|s| s.dispatch_set_parameter(v))
618                            .map_err(|e| e.into_raw()),
619                    ),
620                fidl_fuchsia_net_dhcp::Server_Request::ListOptions { responder: r } => r.send(
621                    server.borrow().dispatch_list_options().as_deref().map_err(|e| e.into_raw()),
622                ),
623                fidl_fuchsia_net_dhcp::Server_Request::ListParameters { responder: r } => r.send(
624                    server.borrow().dispatch_list_parameters().as_deref().map_err(|e| e.into_raw()),
625                ),
626                fidl_fuchsia_net_dhcp::Server_Request::ResetOptions { responder: r } => {
627                    r.send(server.borrow_mut().dispatch_reset_options().map_err(|e| e.into_raw()))
628                }
629                fidl_fuchsia_net_dhcp::Server_Request::ResetParameters { responder: r } => r.send(
630                    server
631                        .borrow_mut()
632                        .if_disabled(|s| s.dispatch_reset_parameters(&default_params))
633                        .map_err(|e| e.into_raw()),
634                ),
635                fidl_fuchsia_net_dhcp::Server_Request::ClearLeases { responder: r } => r.send(
636                    server.borrow_mut().dispatch_clear_leases().map_err(zx::Status::into_raw),
637                ),
638            }
639            .map(|()| socket_sink)
640        })
641        .await
642        // Discard the socket sink.
643        .map(|_socket_sink: C| ())
644}
645
646#[cfg(test)]
647mod tests {
648    use super::*;
649    use dhcpv4::configuration::ServerParameters;
650    use futures::FutureExt;
651    use futures::sink::drain;
652    use net_declare::{fidl_ip_v4, std_ip_v4};
653
654    #[derive(Debug, Eq, PartialEq)]
655    struct CannedSocket {
656        name: String,
657        src: Ipv4Addr,
658    }
659
660    struct CannedDispatcher {
661        params: Option<ServerParameters>,
662        mock_leases: u32,
663    }
664
665    impl CannedDispatcher {
666        fn new() -> Self {
667            Self { params: None, mock_leases: 0 }
668        }
669    }
670
671    impl SocketServerDispatcher for CannedDispatcher {
672        type Socket = CannedSocket;
673
674        fn create_socket(name: &str, src: Ipv4Addr) -> std::io::Result<Self::Socket> {
675            let name = name.to_string();
676            Ok(CannedSocket { name, src })
677        }
678
679        fn dispatch_message(&mut self, mut msg: Message) -> Result<ServerAction, ServerError> {
680            msg.op = dhcpv4::protocol::OpCode::BOOTREPLY;
681            Ok(ServerAction::SendResponse(msg, ResponseTarget::Broadcast))
682        }
683
684        fn create_sockets(
685            params: &configuration::ServerParameters,
686        ) -> std::io::Result<Vec<SocketWithId<Self::Socket>>> {
687            let configuration::ServerParameters { bound_device_names, .. } = params;
688            bound_device_names
689                .iter()
690                .map(String::as_str)
691                .enumerate()
692                .map(|(iface_id, name)| {
693                    let iface_id = std::convert::TryInto::try_into(iface_id).map_err(|e| {
694                        std::io::Error::new(
695                            std::io::ErrorKind::InvalidInput,
696                            format!("interface id {} out of range: {}", iface_id, e),
697                        )
698                    })?;
699                    let socket = Self::create_socket(name, Ipv4Addr::UNSPECIFIED)?;
700                    Ok(SocketWithId { socket, iface_id })
701                })
702                .collect()
703        }
704    }
705
706    impl ServerDispatcher for CannedDispatcher {
707        fn try_validate_parameters(&self) -> Result<&ServerParameters, zx::Status> {
708            self.params.as_ref().ok_or(zx::Status::INVALID_ARGS)
709        }
710
711        fn dispatch_get_option(
712            &self,
713            _code: fidl_fuchsia_net_dhcp::OptionCode,
714        ) -> Result<fidl_fuchsia_net_dhcp::Option_, zx::Status> {
715            Ok(fidl_fuchsia_net_dhcp::Option_::SubnetMask(fidl_ip_v4!("0.0.0.0")))
716        }
717        fn dispatch_get_parameter(
718            &self,
719            _name: fidl_fuchsia_net_dhcp::ParameterName,
720        ) -> Result<fidl_fuchsia_net_dhcp::Parameter, zx::Status> {
721            Ok(fidl_fuchsia_net_dhcp::Parameter::Lease(fidl_fuchsia_net_dhcp::LeaseLength {
722                default: None,
723                max: None,
724                ..Default::default()
725            }))
726        }
727        fn dispatch_set_option(
728            &mut self,
729            _value: fidl_fuchsia_net_dhcp::Option_,
730        ) -> Result<(), zx::Status> {
731            Ok(())
732        }
733        fn dispatch_set_parameter(
734            &mut self,
735            _value: fidl_fuchsia_net_dhcp::Parameter,
736        ) -> Result<(), zx::Status> {
737            Ok(())
738        }
739        fn dispatch_list_options(&self) -> Result<Vec<fidl_fuchsia_net_dhcp::Option_>, zx::Status> {
740            Ok(vec![])
741        }
742        fn dispatch_list_parameters(
743            &self,
744        ) -> Result<Vec<fidl_fuchsia_net_dhcp::Parameter>, zx::Status> {
745            Ok(vec![])
746        }
747        fn dispatch_reset_options(&mut self) -> Result<(), zx::Status> {
748            Ok(())
749        }
750        fn dispatch_reset_parameters(
751            &mut self,
752            _defaults: &dhcpv4::configuration::ServerParameters,
753        ) -> Result<(), zx::Status> {
754            Ok(())
755        }
756        fn dispatch_clear_leases(&mut self) -> Result<(), zx::Status> {
757            self.mock_leases = 0;
758            Ok(())
759        }
760    }
761
762    const DEFAULT_DEVICE_NAME: &str = "foo13";
763
764    fn default_params() -> dhcpv4::configuration::ServerParameters {
765        dhcpv4::configuration::ServerParameters {
766            server_ips: vec![std_ip_v4!("192.168.0.1")],
767            lease_length: dhcpv4::configuration::LeaseLength {
768                default_seconds: 86400,
769                max_seconds: 86400,
770            },
771            managed_addrs: dhcpv4::configuration::ManagedAddresses {
772                mask: dhcpv4::configuration::SubnetMask::new(prefix_length_v4!(25)),
773                pool_range_start: std_ip_v4!("192.168.0.0"),
774                pool_range_stop: std_ip_v4!("192.168.0.0"),
775            },
776            permitted_macs: dhcpv4::configuration::PermittedMacs(vec![]),
777            static_assignments: dhcpv4::configuration::StaticAssignments(HashMap::new()),
778            arp_probe: false,
779            bound_device_names: vec![DEFAULT_DEVICE_NAME.to_string()],
780        }
781    }
782
783    async fn run_with_server<T, F, Fut>(f: F) -> T
784    where
785        F: Fn(fidl_fuchsia_net_dhcp::Server_Proxy) -> Fut,
786        Fut: Future<Output = T>,
787    {
788        let (proxy, stream) =
789            fidl::endpoints::create_proxy_and_stream::<fidl_fuchsia_net_dhcp::Server_Marker>();
790        let server = RefCell::new(ServerDispatcherRuntime::new(CannedDispatcher::new()));
791
792        let defaults = default_params();
793        futures::select! {
794            res = f(proxy).fuse() => res,
795            res = run_server(stream, &server, &defaults, drain()).fuse() => {
796                unreachable!("server finished before request: {:?}", res)
797            },
798        }
799    }
800
801    #[fuchsia::test(logging = false)]
802    async fn get_option_with_subnet_mask_returns_subnet_mask() {
803        run_with_server(|proxy| async move {
804            assert_eq!(
805                proxy
806                    .get_option(fidl_fuchsia_net_dhcp::OptionCode::SubnetMask)
807                    .await
808                    .expect("get_option failed"),
809                Ok(fidl_fuchsia_net_dhcp::Option_::SubnetMask(fidl_ip_v4!("0.0.0.0")))
810            );
811        })
812        .await
813    }
814
815    #[fuchsia::test(allow_stalls = false)]
816    async fn get_parameter_with_lease_length_returns_lease_length() {
817        run_with_server(|proxy| async move {
818            assert_eq!(
819                proxy
820                    .get_parameter(fidl_fuchsia_net_dhcp::ParameterName::LeaseLength)
821                    .await
822                    .expect("get_parameter failed"),
823                Ok(fidl_fuchsia_net_dhcp::Parameter::Lease(fidl_fuchsia_net_dhcp::LeaseLength {
824                    default: None,
825                    max: None,
826                    ..Default::default()
827                }))
828            );
829        })
830        .await
831    }
832
833    #[fuchsia::test(logging = false)]
834    async fn set_option_with_subnet_mask_returns_unit() {
835        run_with_server(|proxy| async move {
836            assert_eq!(
837                proxy
838                    .set_option(&fidl_fuchsia_net_dhcp::Option_::SubnetMask(fidl_ip_v4!("0.0.0.0")))
839                    .await
840                    .expect("set_option failed"),
841                Ok(())
842            );
843        })
844        .await
845    }
846
847    #[fuchsia::test(logging = false)]
848    async fn set_parameter_with_lease_length_returns_unit() {
849        run_with_server(|proxy| async move {
850            assert_eq!(
851                proxy
852                    .set_parameter(&fidl_fuchsia_net_dhcp::Parameter::Lease(
853                        fidl_fuchsia_net_dhcp::LeaseLength {
854                            default: None,
855                            max: None,
856                            ..Default::default()
857                        },
858                    ))
859                    .await
860                    .expect("set_parameter failed"),
861                Ok(())
862            );
863        })
864        .await
865    }
866
867    #[fuchsia::test(logging = false)]
868    async fn list_options_returns_empty_vec() {
869        run_with_server(|proxy| async move {
870            assert_eq!(proxy.list_options().await.expect("list_options failed"), Ok(Vec::new()));
871        })
872        .await
873    }
874
875    #[fuchsia::test(logging = false)]
876    async fn list_parameters_returns_empty_vec() {
877        run_with_server(|proxy| async move {
878            assert_eq!(
879                proxy.list_parameters().await.expect("list_parameters failed"),
880                Ok(Vec::new())
881            );
882        })
883        .await
884    }
885
886    #[fuchsia::test(logging = false)]
887    async fn reset_options_returns_unit() {
888        run_with_server(|proxy| async move {
889            assert_eq!(proxy.reset_options().await.expect("reset_options failed"), Ok(()));
890        })
891        .await
892    }
893
894    #[fuchsia::test(logging = false)]
895    async fn reset_parameters_returns_unit() {
896        run_with_server(|proxy| async move {
897            assert_eq!(proxy.reset_parameters().await.expect("reset_parameters failed"), Ok(()));
898        })
899        .await
900    }
901
902    #[fuchsia::test(logging = false)]
903    async fn clear_leases_returns_unit() {
904        run_with_server(|proxy| async move {
905            assert_eq!(proxy.clear_leases().await.expect("clear_leases failed"), Ok(()));
906        })
907        .await
908    }
909
910    #[fuchsia::test(logging = false)]
911    async fn start_stop_server() {
912        let (proxy, stream) =
913            fidl::endpoints::create_proxy_and_stream::<fidl_fuchsia_net_dhcp::Server_Marker>();
914        let (socket_sink, mut socket_stream) =
915            futures::channel::mpsc::channel::<ServerSocketCollection<CannedSocket>>(1);
916
917        let server = RefCell::new(ServerDispatcherRuntime::new(CannedDispatcher::new()));
918        // Set default parameters to the server so we can create sockets.
919        server.borrow_mut().params = Some(default_params());
920
921        let defaults = default_params();
922
923        // Set mock leases that should not change when the server is disabled.
924        server.borrow_mut().mock_leases = 1;
925
926        let test_fut = async {
927            for () in std::iter::repeat(()).take(3) {
928                assert!(
929                    !proxy.is_serving().await.expect("query server status request"),
930                    "server should not be serving"
931                );
932
933                proxy
934                    .start_serving()
935                    .await
936                    .expect("start_serving failed")
937                    .map_err(zx::Status::err_from_raw)
938                    .expect("start_serving returned an error");
939
940                let ServerSocketCollection { sockets, abort_registration } =
941                    socket_stream.next().await.expect("Socket stream ended unexpectedly");
942
943                // Assert that the sockets that would be created are correct.
944                assert_eq!(
945                    sockets,
946                    vec![SocketWithId {
947                        socket: CannedSocket {
948                            name: DEFAULT_DEVICE_NAME.to_string(),
949                            src: Ipv4Addr::UNSPECIFIED
950                        },
951                        iface_id: 0
952                    }]
953                );
954
955                // Create a dummy future that should be aborted when we disable the
956                // server.
957                let dummy_fut = futures::future::Abortable::new(
958                    futures::future::pending::<()>(),
959                    abort_registration,
960                );
961
962                assert!(
963                    proxy.is_serving().await.expect("query server status request"),
964                    "server should be serving"
965                );
966
967                proxy.stop_serving().await.expect("stop_serving failed");
968
969                // Dummy future was aborted.
970                assert_eq!(dummy_fut.await, Err(futures::future::Aborted {}));
971                // Leases were not cleared.
972                assert_eq!(server.borrow().mock_leases, 1);
973
974                assert!(
975                    !proxy.is_serving().await.expect("query server status request"),
976                    "server should no longer be serving"
977                );
978            }
979        };
980
981        futures::select! {
982            res = test_fut.fuse() => res,
983            res = run_server(stream, &server, &defaults, socket_sink).fuse() => {
984                unreachable!("server finished before request: {:?}", res)
985            },
986        };
987    }
988
989    #[fuchsia::test(logging = false)]
990    async fn start_server_fails_on_bad_params() {
991        let (proxy, stream) =
992            fidl::endpoints::create_proxy_and_stream::<fidl_fuchsia_net_dhcp::Server_Marker>();
993        let server = RefCell::new(ServerDispatcherRuntime::new(CannedDispatcher::new()));
994
995        let defaults = default_params();
996        let res = futures::select! {
997            res = proxy.start_serving().fuse() => res.expect("start_serving failed"),
998            res = run_server(stream, &server, &defaults, drain()).fuse() => {
999                unreachable!("server finished before request: {:?}", res)
1000            },
1001        }
1002        .map_err(zx::Status::err_from_raw);
1003
1004        // Must have failed to start the server.
1005        assert_eq!(res, Err(zx::Status::INVALID_ARGS));
1006        // No abort handler must've been set.
1007        assert!(server.borrow().abort_handle.is_none());
1008    }
1009
1010    #[fuchsia::test(logging = false)]
1011    async fn start_server_fails_on_missing_interface_names() {
1012        let (proxy, stream) =
1013            fidl::endpoints::create_proxy_and_stream::<fidl_fuchsia_net_dhcp::Server_Marker>();
1014        let server = RefCell::new(ServerDispatcherRuntime::new(CannedDispatcher::new()));
1015
1016        let defaults = dhcpv4::configuration::ServerParameters {
1017            bound_device_names: Vec::new(),
1018            ..default_params()
1019        };
1020        server.borrow_mut().params = Some(defaults.clone());
1021
1022        let res = futures::select! {
1023            res = proxy.start_serving().fuse() => res.expect("start_serving failed"),
1024            res = run_server(stream, &server, &defaults, drain()).fuse() => {
1025                unreachable!("server finished before request: {:?}", res)
1026            },
1027        }
1028        .map_err(zx::Status::err_from_raw);
1029
1030        // Must have failed to start the server.
1031        assert_eq!(res, Err(zx::Status::INVALID_ARGS));
1032        // No abort handler must've been set.
1033        assert!(server.borrow().abort_handle.is_none());
1034    }
1035
1036    #[fuchsia::test(logging = false)]
1037    async fn disallow_change_parameters_if_enabled() {
1038        let (proxy, stream) =
1039            fidl::endpoints::create_proxy_and_stream::<fidl_fuchsia_net_dhcp::Server_Marker>();
1040
1041        let server = RefCell::new(ServerDispatcherRuntime::new(CannedDispatcher::new()));
1042        // Set default parameters to the server so we can create sockets.
1043        server.borrow_mut().params = Some(default_params());
1044
1045        let defaults = default_params();
1046
1047        let test_fut = async {
1048            proxy
1049                .start_serving()
1050                .await
1051                .expect("start_serving failed")
1052                .map_err(zx::Status::err_from_raw)
1053                .expect("start_serving returned an error");
1054
1055            // SetParameter disallowed when the server is enabled.
1056            assert_eq!(
1057                proxy
1058                    .set_parameter(&fidl_fuchsia_net_dhcp::Parameter::Lease(
1059                        fidl_fuchsia_net_dhcp::LeaseLength {
1060                            default: None,
1061                            max: None,
1062                            ..Default::default()
1063                        },
1064                    ))
1065                    .await
1066                    .expect("set_parameter FIDL failure")
1067                    .map_err(zx::Status::err_from_raw),
1068                Err(zx::Status::BAD_STATE)
1069            );
1070
1071            // ResetParameters disallowed when the server is enabled.
1072            assert_eq!(
1073                proxy
1074                    .reset_parameters()
1075                    .await
1076                    .expect("reset_parameters FIDL failure")
1077                    .map_err(zx::Status::err_from_raw),
1078                Err(zx::Status::BAD_STATE)
1079            );
1080        };
1081
1082        futures::select! {
1083            res = test_fut.fuse() => res,
1084            res = run_server(stream, &server, &defaults, drain()).fuse() => {
1085                unreachable!("server finished before request: {:?}", res)
1086            },
1087        };
1088    }
1089
1090    /// Test that a malformed message does not cause MessageHandler to return an
1091    /// error.
1092    #[test]
1093    fn handle_failed_parse() {
1094        let server = RefCell::new(ServerDispatcherRuntime::new(CannedDispatcher::new()));
1095        let mut handler = MessageHandler::new(&server);
1096        assert_matches::assert_matches!(
1097            handler.handle_from_sender(
1098                &[0xFF, 0x00, 0xBA, 0x03],
1099                std::net::SocketAddrV4::new(Ipv4Addr::UNSPECIFIED.into(), 0),
1100            ),
1101            Ok(None)
1102        );
1103    }
1104}