Skip to main content

bt_obex/server/
mod.rs

1// Copyright 2023 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 fuchsia_bluetooth::types::Channel;
6use futures::future::Future;
7use futures::stream::StreamExt;
8use log::{info, trace, warn};
9use packet_encoding::{Decodable, Encodable};
10
11use crate::error::{Error, PacketError};
12use crate::header::{
13    ConnectionIdentifier, Header, HeaderIdentifier, HeaderSet, SingleResponseMode,
14};
15use crate::operation::{OpCode, RequestPacket, ResponseCode, ResponsePacket, SetPathFlags};
16pub use crate::transport::TransportType;
17use crate::transport::max_packet_size_from_transport;
18use futures::SinkExt;
19
20/// Defines an interface for handling OBEX requests. All profiles & services should implement this
21/// interface.
22mod handler;
23pub use handler::{ObexOperationError, ObexServerHandler, new_operation_error};
24
25/// Implements the OBEX GET operation.
26mod get;
27use get::GetOperation;
28
29/// Implements the OBEX PUT operation.
30mod put;
31use put::PutOperation;
32
33/// Represents a request to be handled by the OBEX Server during a multi-step operation.
34#[derive(Debug)]
35pub enum OperationRequest {
36    /// Request to send response packets to the remote peer.
37    SendPackets(Vec<ResponsePacket>),
38    /// Request to get informational headers describing a payload from the upper layer application
39    /// -- occurs in a GET operation.
40    GetApplicationInfo(HeaderSet),
41    /// Request to get the payload from the upper layer application -- occurs in a GET operation.
42    GetApplicationData(HeaderSet),
43    /// Request to give the payload to the upper layer application -- occurs in a PUT operation.
44    PutApplicationData(Vec<u8>, HeaderSet),
45    /// Request to delete the payload in the upper layer application -- occurs in a PUT operation.
46    DeleteApplicationData(HeaderSet),
47    /// No action needed.
48    None,
49}
50
51impl OperationRequest {
52    pub fn single_response(packet: ResponsePacket) -> Self {
53        Self::SendPackets(vec![packet])
54    }
55}
56
57/// Represents a response from the upper layer application during a multi-step operation.
58#[derive(Debug)]
59pub enum ApplicationResponse {
60    /// The application responded successfully to get GET information request by providing
61    /// informational headers.
62    GetInfo(HeaderSet),
63    /// The application responded successfully to the GET request by providing the data payload
64    /// and informational headers.
65    GetData((Vec<u8>, HeaderSet)),
66    /// The application responded successfully to the PUT request.
67    Put,
68}
69
70impl ApplicationResponse {
71    #[cfg(test)]
72    fn accept_get(data: Vec<u8>, headers: HeaderSet) -> Result<Self, ObexOperationError> {
73        Ok(ApplicationResponse::GetData((data, headers)))
74    }
75
76    #[cfg(test)]
77    fn accept_get_info(headers: HeaderSet) -> Result<Self, ObexOperationError> {
78        Ok(ApplicationResponse::GetInfo(headers))
79    }
80
81    #[cfg(test)]
82    fn accept_put() -> Result<Self, ObexOperationError> {
83        Ok(ApplicationResponse::Put)
84    }
85}
86
87/// An interface for implementing a multi-step OBEX operation. Currently, the only two such
88/// operations are GET and PUT.
89/// See OBEX 1.5 Sections 3.4.3 & 3.4.4.
90pub trait ServerOperation {
91    /// Returns the current SRM mode of the operation.
92    fn srm_status(&self) -> SingleResponseMode;
93
94    /// Checks the provided `headers` for the SRM flag and returns the negotiated SRM mode if
95    /// present, None otherwise.
96    fn check_headers_for_srm(
97        srm_supported_locally: bool,
98        headers: &HeaderSet,
99    ) -> Option<SingleResponseMode>
100    where
101        Self: Sized,
102    {
103        let Some(Header::SingleResponseMode(srm)) =
104            headers.get(&HeaderIdentifier::SingleResponseMode)
105        else {
106            trace!("No SRM header in request");
107            return None;
108        };
109
110        // If both parties support SRM, then it can be enabled.
111        if srm_supported_locally && *srm == SingleResponseMode::Enable {
112            Some(SingleResponseMode::Enable)
113        } else {
114            // Otherwise, either we don't support it locally, or the peer requested to disable it.
115            Some(SingleResponseMode::Disable)
116        }
117    }
118
119    /// Returns true if the operation is complete (e.g. all response packets have been sent).
120    fn is_complete(&self) -> bool;
121
122    /// Handle a `request` packet received from the OBEX client.
123    /// Returns an `OperationRequest` to be handled by the OBEX server on success, Error if the
124    /// request was invalid or couldn't be handled.
125    fn handle_peer_request(&mut self, request: RequestPacket) -> Result<OperationRequest, Error>;
126
127    /// Handle a response received from the upper layer application profile.
128    /// `response` is Ok<T> if the application accepted the GET or PUT request.
129    /// `response` is Err<E> if the application rejected the GET or PUT request.
130    /// Returns response packets to be sent to the remote peer if the application `response` was
131    /// successfully handled.
132    /// Returns Error if there was an internal operation error.
133    fn handle_application_response(
134        &mut self,
135        response: Result<ApplicationResponse, ObexOperationError>,
136    ) -> Result<Vec<ResponsePacket>, Error>;
137}
138
139#[derive(Clone, Copy, Debug, Default, PartialEq)]
140enum ConnectionStatus {
141    /// The transport is created but the CONNECT operation has not been completed.
142    #[default]
143    Initialized,
144    /// The transport is connected and the CONNECT operation has been completed.
145    /// `id` contains the optional identifier for this connection. It is typically set when the
146    /// OBEX client requests a directed OBEX connection to a specific Service by including the
147    /// `Target` header in the CONNECT request.
148    Connected { id: Option<ConnectionIdentifier> },
149    /// The transport is connected but a DISCONNECT request has been received. The `ObexServer`
150    /// will no longer process requests from the remote peer.
151    DisconnectReceived,
152}
153
154impl ConnectionStatus {
155    #[cfg(test)]
156    fn connected_no_id() -> Self {
157        Self::Connected { id: None }
158    }
159}
160
161/// Implements the Server role for the OBEX protocol.
162/// Provides an interface for receiving and responding to OBEX requests made by a remote OBEX client
163/// service. Supports the operations defined in OBEX 1.5.
164pub struct ObexServer {
165    /// The current connection status of the server.
166    connected: ConnectionStatus,
167    /// The maximum OBEX packet length for this OBEX session.
168    max_packet_size: u16,
169    /// The active OBEX operation. The only two supported multi-step operations are GET and PUT.
170    /// This is Some<T> when an operation is in progress, and None otherwise. There can only be one
171    /// active multi-step operation. An operation is considered complete when
172    /// `ServerOperation::is_complete` returns true.
173    /// The active operation is cleaned up lazily -- when a request to start a new operation is
174    /// received, the previously finished operation is removed.
175    active_operation: Option<Box<dyn ServerOperation>>,
176    /// The data channel that is used to read & write OBEX packets.
177    channel: Channel,
178    /// The type of transport used for the OBEX connection (RFCOMM or L2CAP).
179    type_: TransportType,
180    /// The handler provided by the application profile. This handler should implement the
181    /// operations defined in OBEX 1.5 and will be used to provide a response to an incoming
182    /// request made by the remote OBEX client.
183    handler: Box<dyn ObexServerHandler>,
184}
185
186impl ObexServer {
187    /// The default Connection Identifier used for directed OBEX connections.
188    /// Because a single `ObexServer` services a single transport (L2CAP or RFCOMM), this
189    /// identifier does not multiplex anything and is only sent in the CONNECT response.
190    /// This value is arbitrarily chosen and will be included in all subsequent requests made by
191    /// the remote OBEX Client.
192    const DIRECTED_CONNECTION_ID: ConnectionIdentifier = ConnectionIdentifier(1);
193
194    pub fn new(
195        channel: Channel,
196        type_: TransportType,
197        handler: Box<dyn ObexServerHandler>,
198    ) -> Self {
199        let max_packet_size = max_packet_size_from_transport(channel.max_tx_size());
200        Self {
201            connected: ConnectionStatus::default(),
202            max_packet_size,
203            active_operation: None,
204            channel,
205            type_,
206            handler,
207        }
208    }
209
210    /// Returns `true` if the OBEX connection is currently active (e.g. CONNECT operation done).
211    fn is_connected(&self) -> bool {
212        matches!(self.connected, ConnectionStatus::Connected { .. })
213    }
214
215    fn set_connection_status(&mut self, status: ConnectionStatus) {
216        self.connected = status;
217    }
218
219    fn set_max_packet_size(&mut self, peer_max_packet_size: u16) {
220        // Use the smaller of the peer max and local max for maximum compatibility.
221        let min_ = std::cmp::min(peer_max_packet_size, self.max_packet_size);
222        self.max_packet_size = min_;
223        trace!("Max packet size set to {}", self.max_packet_size);
224    }
225
226    /// Encodes and sends the OBEX `data` to the remote peer.
227    /// Returns Error if the send operation could not be completed.
228    async fn send(&mut self, data: impl Encodable<Error = PacketError>) -> Result<(), Error> {
229        let mut buf = vec![0; data.encoded_len()];
230        if buf.len() > self.max_packet_size as usize {
231            return Err(Error::packet_too_large(buf.len(), self.max_packet_size as usize));
232        }
233
234        data.encode(&mut buf[..])?;
235        self.channel.send(buf).await?;
236        Ok(())
237    }
238
239    async fn connect_request(&mut self, request: RequestPacket) -> Result<ResponsePacket, Error> {
240        // Parse the additional data first - the data length is already validated during decoding.
241        let data = request.data();
242        let version = data[0];
243        let flags = data[1];
244        let peer_max_packet_size = u16::from_be_bytes(data[2..4].try_into().unwrap());
245        trace!(version, flags, peer_max_packet_size; "Additional data in CONNECT request");
246        self.set_max_packet_size(peer_max_packet_size);
247
248        let headers = HeaderSet::from(request);
249
250        // The connection can optionally be considered "directed" if the Client provides a Target
251        // UUID identifying the service.
252        let id = if headers.contains_header(&HeaderIdentifier::Target) {
253            Some(Self::DIRECTED_CONNECTION_ID)
254        } else {
255            None
256        };
257        let (code, response_headers) = match self.handler.connect(headers).await {
258            Ok(mut headers) => {
259                trace!("Application accepted CONNECT request");
260                let _ = headers.try_add_connection_id(&id);
261                self.set_connection_status(ConnectionStatus::Connected { id });
262                (ResponseCode::Ok, headers)
263            }
264            Err(reject_parameters) => {
265                trace!("Application rejected CONNECT request");
266                reject_parameters
267            }
268        };
269        let response_packet =
270            ResponsePacket::new_connect(code, self.max_packet_size, response_headers);
271        Ok(response_packet)
272    }
273
274    /// Handles a Disconnect request made by the remote OBEX client.
275    /// Returns a response packet to be sent on success, Error if the request couldn't be handled.
276    async fn disconnect_request(
277        &mut self,
278        request: RequestPacket,
279    ) -> Result<ResponsePacket, Error> {
280        let headers = HeaderSet::from(request);
281        let response_headers = self.handler.disconnect(headers).await;
282        let response_packet = ResponsePacket::new_disconnect(response_headers);
283        self.set_connection_status(ConnectionStatus::DisconnectReceived);
284        Ok(response_packet)
285    }
286
287    /// Handles a SetPath request made by the remote OBEX client.
288    /// Returns a response packet to be sent on success, Error if the request couldn't be handled.
289    async fn setpath_request(&mut self, request: RequestPacket) -> Result<ResponsePacket, Error> {
290        if !self.is_connected() {
291            return Err(Error::operation(OpCode::SetPath, "CONNECT not completed"));
292        }
293        // Parse the additional data first - the data length is already validated during decoding.
294        // Only the `flags` field is used in OBEX 1.5. `constants` is RFA.
295        let data = request.data();
296        let flags = SetPathFlags::from_bits_truncate(data[0]);
297        let backup = flags.contains(SetPathFlags::BACKUP);
298        let create = !flags.contains(SetPathFlags::DONT_CREATE);
299
300        let headers = HeaderSet::from(request);
301        let (code, response_headers) = match self.handler.set_path(headers, backup, create).await {
302            Ok(headers) => {
303                trace!("Application accepted SETPATH request");
304                (ResponseCode::Ok, headers)
305            }
306            Err(reject_parameters) => {
307                trace!("Application rejected SETPATH request");
308                reject_parameters
309            }
310        };
311        let response_packet = ResponsePacket::new_setpath(code, response_headers);
312        Ok(response_packet)
313    }
314
315    /// Potentially initializes a new multi-step operation.
316    /// Returns true if a new operation was initialized, false otherwise.
317    fn maybe_start_new_operation(&mut self, code: &OpCode) -> bool {
318        if self.active_operation.as_ref().is_some_and(|o| !o.is_complete()) {
319            return false;
320        }
321
322        let op: Box<dyn ServerOperation> = match code {
323            OpCode::Get | OpCode::GetFinal => {
324                Box::new(GetOperation::new(self.max_packet_size, self.type_.srm_supported()))
325            }
326            OpCode::Put | OpCode::PutFinal => {
327                Box::new(PutOperation::new(self.type_.srm_supported()))
328            }
329            _ => unreachable!("only called from `Self::multistep_request`"),
330        };
331        trace!("Started new operation ({code:?})");
332        self.active_operation = Some(op);
333        return true;
334    }
335
336    /// Handles a request made by the remote OBEX client for a potentially multi-step
337    /// operation (PUT or GET).
338    /// Returns response packets to be sent to the peer on success, Error if the request can't
339    /// be handled or is invalid.
340    async fn multistep_request(
341        &mut self,
342        request: RequestPacket,
343    ) -> Result<Vec<ResponsePacket>, Error> {
344        let _ = self.maybe_start_new_operation(request.code());
345        let operation = self.active_operation.as_mut().expect("just initialized");
346
347        let application_response = match operation.handle_peer_request(request) {
348            Ok(OperationRequest::SendPackets(responses)) => return Ok(responses),
349            Ok(OperationRequest::GetApplicationInfo(info_headers)) => {
350                self.handler.get_info(info_headers).await.map(|x| ApplicationResponse::GetInfo(x))
351            }
352            Ok(OperationRequest::GetApplicationData(request_headers)) => self
353                .handler
354                .get_data(request_headers)
355                .await
356                .map(|x| ApplicationResponse::GetData(x)),
357            Ok(OperationRequest::PutApplicationData(data, request_headers)) => {
358                self.handler.put(data, request_headers).await.map(|_| ApplicationResponse::Put)
359            }
360            Ok(OperationRequest::DeleteApplicationData(request_headers)) => {
361                self.handler.delete(request_headers).await.map(|_| ApplicationResponse::Put)
362            }
363            Ok(OperationRequest::None) => return Ok(vec![]),
364            Err(e) => {
365                warn!("Internal error in operation: {e:?}");
366                return Ok(vec![ResponsePacket::new_no_data(
367                    ResponseCode::InternalServerError,
368                    HeaderSet::new(),
369                )]);
370            }
371        };
372
373        operation.handle_application_response(application_response)
374    }
375
376    /// Processes a raw data `packet` received from the remote peer acting as an OBEX client.
377    /// Returns a list of `ResponsePacket`s to be sent to the peer on success, Error if the request
378    /// was invalid or couldn't be handled.
379    async fn receive_packet(&mut self, packet: Vec<u8>) -> Result<Vec<ResponsePacket>, Error> {
380        if packet.len() > self.max_packet_size as usize {
381            warn!(
382                "Received packet size ({}) exceeds negotiated max_packet_size ({})",
383                packet.len(),
384                self.max_packet_size
385            );
386            return Ok(vec![ResponsePacket::new_no_data(
387                ResponseCode::BadRequest,
388                HeaderSet::new(),
389            )]);
390        }
391
392        let decoded = RequestPacket::decode(&packet[..])?;
393        trace!(packet:? = decoded; "Received request from OBEX client");
394        let response = match decoded.code() {
395            OpCode::Connect => self.connect_request(decoded).await?,
396            OpCode::Disconnect => self.disconnect_request(decoded).await?,
397            OpCode::SetPath => self.setpath_request(decoded).await?,
398            OpCode::Put | OpCode::PutFinal | OpCode::Get | OpCode::GetFinal => {
399                return self.multistep_request(decoded).await;
400            }
401            _code => todo!("Support other OBEX requests"),
402        };
403        Ok(vec![response])
404    }
405
406    pub fn run(mut self) -> impl Future<Output = Result<(), Error>> {
407        async move {
408            while let Some(packet) = self.channel.next().await {
409                match packet {
410                    Ok(bytes) => {
411                        let responses = self.receive_packet(bytes).await?;
412                        for response in responses {
413                            self.send(response).await?;
414                        }
415
416                        // The OBEX Client requested to close the OBEX connection.
417                        if self.connected == ConnectionStatus::DisconnectReceived {
418                            trace!("Disconnect request - closing transport");
419                            return Ok(());
420                        }
421                    }
422                    Err(e) => warn!("Error reading data from transport: {e:?}"),
423                }
424            }
425            info!("Peer disconnected transport");
426            Ok(())
427        }
428    }
429}
430
431#[cfg(test)]
432pub(crate) mod test_utils {
433    use super::*;
434
435    #[track_caller]
436    pub fn expect_single_packet(request: OperationRequest) -> ResponsePacket {
437        let OperationRequest::SendPackets(mut packets) = request else {
438            panic!("Expected outgoing packet request, got: {request:?}");
439        };
440        assert_eq!(packets.len(), 1);
441        packets.pop().unwrap()
442    }
443}
444
445#[cfg(test)]
446mod tests {
447    use super::*;
448    use bt_channel_test_support::{Transport, create_test_channels};
449    use test_case::test_case;
450
451    use assert_matches::assert_matches;
452    use async_test_helpers::expect_stream_pending;
453    use async_utils::PollExt;
454    use fuchsia_async as fasync;
455    use std::pin::pin;
456
457    use crate::header::header_set::{expect_body, expect_end_of_body};
458    use crate::server::handler::test_utils::TestApplicationProfile;
459    use crate::transport::test_utils::{expect_response, send_packet};
460
461    /// Returns an ObexServer, a test only object representing an upper layer profile, & the remote
462    /// peer's side of the transport.
463    fn new_obex_server(
464        transport: Transport,
465        srm: bool,
466    ) -> (ObexServer, TestApplicationProfile, Channel) {
467        let (local, remote) = create_test_channels(transport);
468        let app = TestApplicationProfile::new();
469        let type_ = if srm { TransportType::L2cap } else { TransportType::Rfcomm };
470        let obex_server = ObexServer::new(local, type_, Box::new(app.clone()));
471        (obex_server, app, remote)
472    }
473
474    #[test_case(Transport::Socket ; "socket")]
475    #[test_case(Transport::Fidl ; "fidl")]
476    #[fuchsia::test]
477    fn obex_server_terminates_when_channel_closes(transport: Transport) {
478        let mut exec = fasync::TestExecutor::new();
479        let (obex_server, _test_app, remote) = new_obex_server(transport, /*srm=*/ false);
480
481        let server_fut = obex_server.run();
482        let mut server_fut = pin!(server_fut);
483        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server still active");
484
485        drop(remote);
486        let result = exec.run_until_stalled(&mut server_fut).expect("server finished");
487        assert_matches!(result, Ok(_));
488    }
489
490    #[test_case(Transport::Socket ; "socket")]
491    #[test_case(Transport::Fidl ; "fidl")]
492    #[fuchsia::test]
493    fn connect_accepted_by_app_success(transport: Transport) {
494        let mut exec = fasync::TestExecutor::new();
495        let (obex_server, test_app, mut remote) = new_obex_server(transport, /*srm=*/ false);
496        let server_fut = obex_server.run();
497        let mut server_fut = pin!(server_fut);
498        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
499
500        let connect_request = RequestPacket::new_connect(500, HeaderSet::new());
501        send_packet(&mut exec, &mut remote, connect_request);
502
503        // Expect the ObexServer to receive the request, parse it, ask the application, and reply.
504        // Simulate application accepting the request.
505        let headers = HeaderSet::from_header(Header::Description("foo".into()));
506        test_app.set_response(Ok(headers));
507        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
508
509        // Expect the remote peer to receive our CONNECT response. Our response shouldn't contain
510        // a `ConnectionIdentifier` since the request didn't contain a `Target` header.
511        let expectation = |response: ResponsePacket| {
512            assert_eq!(*response.code(), ResponseCode::Ok);
513            assert_eq!(response.data(), &[0x10, 0, 0x01, 0xf4]);
514            assert!(response.headers().contains_header(&HeaderIdentifier::Description));
515            assert!(!response.headers().contains_header(&HeaderIdentifier::ConnectionId));
516        };
517        expect_response(&mut exec, &mut remote, expectation, OpCode::Connect);
518    }
519
520    #[test_case(Transport::Socket ; "socket")]
521    #[test_case(Transport::Fidl ; "fidl")]
522    #[fuchsia::test]
523    fn directed_connect_accepted_by_app_success(transport: Transport) {
524        let mut exec = fasync::TestExecutor::new();
525        let (obex_server, test_app, mut remote) = new_obex_server(transport, /*srm=*/ false);
526        let server_fut = obex_server.run();
527        let mut server_fut = pin!(server_fut);
528        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
529
530        let request_headers = HeaderSet::from_header(Header::Target(vec![5, 6]));
531        let connect_request = RequestPacket::new_connect(500, request_headers);
532        send_packet(&mut exec, &mut remote, connect_request);
533
534        // Expect the ObexServer to receive the request, parse it, ask the application, and reply.
535        // Simulate application accepting the request.
536        let headers = HeaderSet::from_header(Header::name("foo"));
537        test_app.set_response(Ok(headers));
538        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
539
540        // Expect the remote peer to receive our CONNECT response. Our response should contain a
541        // connection identifier because the `request_headers` contains a `Target` header.
542        let expectation = |response: ResponsePacket| {
543            assert_eq!(*response.code(), ResponseCode::Ok);
544            assert_eq!(response.data(), &[0x10, 0, 0x01, 0xf4]);
545            assert!(response.headers().contains_header(&HeaderIdentifier::Name));
546            assert!(response.headers().contains_header(&HeaderIdentifier::ConnectionId));
547        };
548        expect_response(&mut exec, &mut remote, expectation, OpCode::Connect);
549    }
550
551    #[test_case(Transport::Socket ; "socket")]
552    #[test_case(Transport::Fidl ; "fidl")]
553    #[fuchsia::test]
554    fn connect_rejected_by_app_is_ok(transport: Transport) {
555        let mut exec = fasync::TestExecutor::new();
556        let (obex_server, test_app, mut remote) = new_obex_server(transport, /*srm=*/ false);
557        let server_fut = obex_server.run();
558        let mut server_fut = pin!(server_fut);
559        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
560
561        let connect_request = RequestPacket::new_connect(255, HeaderSet::new());
562        send_packet(&mut exec, &mut remote, connect_request);
563
564        // The ObexServer should receive the request and hand it to the profile - profile rejects.
565        test_app.set_response(Err((ResponseCode::Forbidden, HeaderSet::new())));
566        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
567
568        // Expect the remote peer to receive our negative CONNECT response.
569        let expectation = |response: ResponsePacket| {
570            assert_eq!(*response.code(), ResponseCode::Forbidden);
571            assert_eq!(response.data(), &[0x10, 0, 0x00, 0xff]);
572            let headers = HeaderSet::from(response);
573            assert!(headers.is_empty());
574        };
575        expect_response(&mut exec, &mut remote, expectation, OpCode::Connect);
576    }
577
578    #[test_case(Transport::Socket ; "socket")]
579    #[test_case(Transport::Fidl ; "fidl")]
580    #[fuchsia::test]
581    fn invalid_connect_request_is_error(transport: Transport) {
582        let mut exec = fasync::TestExecutor::new();
583        let (obex_server, _test_app, mut remote) = new_obex_server(transport, /*srm=*/ false);
584
585        let server_fut = obex_server.run();
586        let mut server_fut = pin!(server_fut);
587        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server still active");
588
589        // Invalid CONNECT request. Missing the 2 byte max packet size field.
590        let bytes = [0x80, 0x00, 0x05, 0x00, 0x00];
591        let mut fut = remote.send(bytes.to_vec());
592        exec.run_until_stalled(&mut fut).expect("can send data").expect("write success");
593
594        let result = exec.run_until_stalled(&mut server_fut).expect("terminate due to error");
595        assert_matches!(result, Err(Error::Packet(_)));
596    }
597
598    #[test_case(Transport::Socket ; "socket")]
599    #[test_case(Transport::Fidl ; "fidl")]
600    #[fuchsia::test]
601    fn peer_disconnect_request_terminates_server(transport: Transport) {
602        let mut exec = fasync::TestExecutor::new();
603        let (obex_server, test_app, mut remote) = new_obex_server(transport, /*srm=*/ false);
604        let server_fut = obex_server.run();
605        let mut server_fut = pin!(server_fut);
606        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
607
608        let headers = HeaderSet::from_header(Header::Description("done".into()));
609        let disconnect_request = RequestPacket::new_disconnect(headers);
610        send_packet(&mut exec, &mut remote, disconnect_request);
611
612        // Expect the ObexServer to receive the request, parse it, get response headers from the
613        // application, and reply. Because this is a Disconnect request, the server run loop
614        // should finish.
615        let headers = HeaderSet::from_header(Header::Description("disconnecting".into()));
616        test_app.set_response(Ok(headers));
617        let result =
618            exec.run_until_stalled(&mut server_fut).expect("server terminated from disconnect");
619        assert_matches!(result, Ok(_));
620
621        // Expect the remote peer to receive our DISCONNECT response.
622        let expectation = |response: ResponsePacket| {
623            assert_eq!(*response.code(), ResponseCode::Ok);
624            let headers = HeaderSet::from(response);
625            assert!(headers.contains_header(&HeaderIdentifier::Description));
626        };
627        expect_response(&mut exec, &mut remote, expectation, OpCode::Disconnect);
628    }
629
630    #[test_case(Transport::Socket ; "socket")]
631    #[test_case(Transport::Fidl ; "fidl")]
632    #[fuchsia::test]
633    fn setpath_request_accepted_by_app_success(transport: Transport) {
634        let mut exec = fasync::TestExecutor::new();
635        let (mut obex_server, test_app, mut remote) =
636            new_obex_server(transport, /*srm=*/ false);
637        // Set to the Connected state to bypass CONNECT operation.
638        obex_server.set_connection_status(ConnectionStatus::connected_no_id());
639        let server_fut = obex_server.run();
640        let mut server_fut = pin!(server_fut);
641        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
642
643        let headers = HeaderSet::from_header(Header::name("folder1"));
644        let setpath_request =
645            RequestPacket::new_set_path(SetPathFlags::all(), headers).expect("valid request");
646        send_packet(&mut exec, &mut remote, setpath_request);
647
648        // The ObexServer should receive the request and hand it to the profile - profile accepts.
649        test_app.set_response(Ok(HeaderSet::new()));
650        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
651
652        // Expect the remote peer to receive our SETPATH response.
653        let expectation = |response: ResponsePacket| {
654            assert_eq!(*response.code(), ResponseCode::Ok);
655            let headers = HeaderSet::from(response);
656            assert!(headers.is_empty());
657        };
658        expect_response(&mut exec, &mut remote, expectation, OpCode::SetPath);
659    }
660
661    #[test_case(Transport::Socket ; "socket")]
662    #[test_case(Transport::Fidl ; "fidl")]
663    #[fuchsia::test]
664    fn setpath_request_rejected_by_app_success(transport: Transport) {
665        let mut exec = fasync::TestExecutor::new();
666        let (mut obex_server, test_app, mut remote) =
667            new_obex_server(transport, /*srm=*/ false);
668        // Set to the Connected state to bypass CONNECT operation.
669        obex_server.set_connection_status(ConnectionStatus::connected_no_id());
670        let server_fut = obex_server.run();
671        let mut server_fut = pin!(server_fut);
672        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
673
674        let setpath_request = RequestPacket::new_set_path(SetPathFlags::BACKUP, HeaderSet::new())
675            .expect("valid request");
676        send_packet(&mut exec, &mut remote, setpath_request);
677
678        // The ObexServer should receive the request and hand it to the profile - profile rejects.
679        test_app.set_response(Err((ResponseCode::Forbidden, HeaderSet::new())));
680        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
681
682        // Expect the remote peer to receive our negative SETPATH response.
683        let expectation = |response: ResponsePacket| {
684            assert_eq!(*response.code(), ResponseCode::Forbidden);
685            let headers = HeaderSet::from(response);
686            assert!(headers.is_empty());
687        };
688        expect_response(&mut exec, &mut remote, expectation, OpCode::SetPath);
689    }
690
691    #[test_case(Transport::Socket ; "socket")]
692    #[test_case(Transport::Fidl ; "fidl")]
693    #[fuchsia::test]
694    fn setpath_request_before_connect_is_error(transport: Transport) {
695        let mut exec = fasync::TestExecutor::new();
696        let (obex_server, _test_app, mut remote) = new_obex_server(transport, /*srm=*/ false);
697        let server_fut = obex_server.run();
698        let mut server_fut = pin!(server_fut);
699        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
700
701        let setpath_request = RequestPacket::new_set_path(SetPathFlags::BACKUP, HeaderSet::new())
702            .expect("valid request");
703        send_packet(&mut exec, &mut remote, setpath_request);
704        let result = exec
705            .run_until_stalled(&mut server_fut)
706            .expect("server terminated from invalid setpath");
707        assert_matches!(result, Err(Error::OperationError { operation: OpCode::SetPath, .. }));
708    }
709
710    #[test_case(Transport::Socket ; "socket")]
711    #[test_case(Transport::Fidl ; "fidl")]
712    #[fuchsia::test]
713    fn get_request_accepted_by_app_success(transport: Transport) {
714        let mut exec = fasync::TestExecutor::new();
715        let (mut obex_server, test_app, mut remote) =
716            new_obex_server(transport, /*srm=*/ false);
717        // Set to the Connected state to bypass CONNECT operation.
718        obex_server.set_connection_status(ConnectionStatus::connected_no_id());
719        let server_fut = obex_server.run();
720        let mut server_fut = pin!(server_fut);
721        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
722
723        // Remote asks for information about the payload. The ObexServer should receive the request
724        // and ask the application for the response.
725        let get_request1 =
726            RequestPacket::new_get(HeaderSet::from_header(Header::name("random object")));
727        send_packet(&mut exec, &mut remote, get_request1);
728        // Simulate the application responding with the size of the object.
729        test_app.set_response(Ok(HeaderSet::from_header(Header::Length(0x10))));
730        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
731
732        // The remote peer should receive the info response.
733        let expectation = |response: ResponsePacket| {
734            assert_eq!(*response.code(), ResponseCode::Continue);
735            assert!(response.headers().contains_header(&HeaderIdentifier::Length));
736        };
737        expect_response(&mut exec, &mut remote, expectation, OpCode::Get);
738
739        // Remote sends a GET_FINAL request indicating that it is ready to receive the data payload.
740        let get_request2 = RequestPacket::new_get_final(HeaderSet::new());
741        send_packet(&mut exec, &mut remote, get_request2);
742
743        // The ObexServer should receive the request and hand it to the profile. Set the profile
744        // handler to return a static buffer.
745        let application_response_buf = vec![1, 2, 3, 4, 5, 6];
746        let response_headers = HeaderSet::from_header(Header::Description("foo".into()));
747        test_app.set_get_response((application_response_buf.clone(), response_headers));
748        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
749
750        let expectation = |response: ResponsePacket| {
751            assert_eq!(*response.code(), ResponseCode::Ok);
752            let mut headers = HeaderSet::from(response);
753            assert!(headers.contains_header(&HeaderIdentifier::Description));
754            let received_body = headers.remove_body(/*final_=*/ true).expect("contains body");
755            assert_eq!(received_body, application_response_buf);
756        };
757        expect_response(&mut exec, &mut remote, expectation, OpCode::GetFinal);
758    }
759
760    #[test_case(Transport::Socket ; "socket")]
761    #[test_case(Transport::Fidl ; "fidl")]
762    #[fuchsia::test]
763    fn get_request_rejected_by_app_success(transport: Transport) {
764        let mut exec = fasync::TestExecutor::new();
765        let (mut obex_server, _test_app, mut remote) =
766            new_obex_server(transport, /*srm=*/ false);
767        obex_server.set_connection_status(ConnectionStatus::connected_no_id());
768        let server_fut = obex_server.run();
769        let mut server_fut = pin!(server_fut);
770        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
771
772        // Send an example GET_FINAL request with a header describing the name of the object.
773        let headers = HeaderSet::from_header(Header::name("random object123"));
774        let get_request = RequestPacket::new_get_final(headers);
775        send_packet(&mut exec, &mut remote, get_request);
776
777        // The ObexServer receives request and hands to application. By default, it rejects the
778        // request.
779        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
780        // Expect the peer to received the rejection code.
781        let expectation = |response: ResponsePacket| {
782            assert_eq!(*response.code(), ResponseCode::NotImplemented);
783            assert!(response.headers().is_empty());
784        };
785        expect_response(&mut exec, &mut remote, expectation, OpCode::GetFinal);
786    }
787
788    #[test_case(Transport::Socket ; "socket")]
789    #[test_case(Transport::Fidl ; "fidl")]
790    #[fuchsia::test]
791    fn get_request_with_srm_enabled_success(transport: Transport) {
792        let mut exec = fasync::TestExecutor::new();
793        let (mut obex_server, test_app, mut remote) =
794            new_obex_server(transport, /*srm=*/ true);
795        obex_server.set_connection_status(ConnectionStatus::connected_no_id());
796        obex_server.set_max_packet_size(20); // Set max to something small.
797        let server_fut = obex_server.run();
798        let mut server_fut = pin!(server_fut);
799        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
800
801        // First request asks for information & SRM. Server should receive the request, ask the
802        // application, and negotiate SRM.
803        let headers1 =
804            HeaderSet::from_headers(vec![Header::name("a"), SingleResponseMode::Enable.into()])
805                .unwrap();
806        let get_request1 = RequestPacket::new_get(headers1);
807        send_packet(&mut exec, &mut remote, get_request1);
808        test_app.set_response(Ok(HeaderSet::from_header(Header::Length(0x20))));
809        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
810        // Expect the response packet to the peer to negotiate SRM and contain the application
811        // response.
812        let expectation1 = |response: ResponsePacket| {
813            assert_eq!(*response.code(), ResponseCode::Continue);
814            let Header::SingleResponseMode(SingleResponseMode::Enable) =
815                response.headers().get(&HeaderIdentifier::SingleResponseMode).unwrap()
816            else {
817                panic!("Expected SRM enable in response");
818            };
819            assert!(response.headers().contains_header(&HeaderIdentifier::Length));
820        };
821        expect_response(&mut exec, &mut remote, expectation1, OpCode::Get);
822        // At this point, SRM is considered active for the operation.
823
824        // Second (final) request to get the payload.
825        let get_request2 = RequestPacket::new_get_final(HeaderSet::new());
826        send_packet(&mut exec, &mut remote, get_request2);
827        // The ObexServer should receive the request and hand it to the profile. Set the profile
828        // handler to return a static buffer that must be split across multiple payloads.
829        let application_response_buf = (0..50).collect::<Vec<u8>>();
830        test_app.set_get_response((application_response_buf, HeaderSet::new()));
831        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
832
833        // Since SRM is enabled, we expect consecutive packets containing the payload - no requests
834        // made by the remote in between.
835        let expected_bufs = vec![
836            (0..14).collect::<Vec<u8>>(),
837            (14..28).collect::<Vec<u8>>(),
838            (28..42).collect::<Vec<u8>>(),
839        ];
840        for expected_buf in expected_bufs {
841            let expectation = |response: ResponsePacket| {
842                assert_eq!(*response.code(), ResponseCode::Continue);
843                expect_body(response.headers(), expected_buf);
844            };
845            expect_response(&mut exec, &mut remote, expectation, OpCode::Get);
846        }
847
848        // Final packet has remaining bytes and operation is complete.
849        let final_expectation = |response: ResponsePacket| {
850            assert_eq!(*response.code(), ResponseCode::Ok);
851            expect_end_of_body(response.headers(), (42..50).collect::<Vec<u8>>());
852        };
853        expect_response(&mut exec, &mut remote, final_expectation, OpCode::GetFinal);
854    }
855
856    #[test_case(Transport::Socket ; "socket")]
857    #[test_case(Transport::Fidl ; "fidl")]
858    #[fuchsia::test]
859    fn put_request_accepted_by_app_success(transport: Transport) {
860        let mut exec = fasync::TestExecutor::new();
861        let (mut obex_server, test_app, mut remote) =
862            new_obex_server(transport, /*srm=*/ false);
863        // Set to the Connected state to bypass CONNECT operation.
864        obex_server.set_connection_status(ConnectionStatus::connected_no_id());
865        let server_fut = obex_server.run();
866        let mut server_fut = pin!(server_fut);
867        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
868
869        let headers = HeaderSet::from_headers(vec![
870            Header::name("random object"),
871            Header::EndOfBody(vec![1, 2, 3, 4, 5]),
872        ])
873        .unwrap();
874        let put_request = RequestPacket::new_put_final(headers);
875        send_packet(&mut exec, &mut remote, put_request);
876
877        // The ObexServer should receive the request and hand it to the profile. Profile accepts.
878        test_app.set_put_response(Ok(()));
879        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
880        // Verify profile received correct data.
881        let (rec_data, rec_headers) = test_app.put_data();
882        assert_eq!(rec_data, vec![1, 2, 3, 4, 5]);
883        assert!(rec_headers.contains_header(&HeaderIdentifier::Name));
884
885        let expectation = |response: ResponsePacket| {
886            assert_eq!(*response.code(), ResponseCode::Ok);
887            assert!(response.headers().is_empty());
888        };
889        expect_response(&mut exec, &mut remote, expectation, OpCode::PutFinal);
890    }
891
892    #[test_case(Transport::Socket ; "socket")]
893    #[test_case(Transport::Fidl ; "fidl")]
894    #[fuchsia::test]
895    fn put_request_with_srm_enabled_success(transport: Transport) {
896        let mut exec = fasync::TestExecutor::new();
897        let (mut obex_server, test_app, mut remote) =
898            new_obex_server(transport, /*srm=*/ true);
899        obex_server.set_connection_status(ConnectionStatus::connected_no_id());
900        let server_fut = obex_server.run();
901        let mut server_fut = pin!(server_fut);
902        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
903
904        // First request asks to enable SRM and provides some info.
905        let headers1 = HeaderSet::from_headers(vec![
906            Header::name("my file"),
907            SingleResponseMode::Enable.into(),
908        ])
909        .unwrap();
910        let put_request1 = RequestPacket::new_put(headers1);
911        send_packet(&mut exec, &mut remote, put_request1);
912        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
913        // Expect the Obex Server to positively respond, and enable SRM. Subsequent requests won't
914        // receive a response.
915        let expectation1 = |response: ResponsePacket| {
916            assert_eq!(*response.code(), ResponseCode::Continue);
917            let Header::SingleResponseMode(SingleResponseMode::Enable) =
918                response.headers().get(&HeaderIdentifier::SingleResponseMode).unwrap()
919            else {
920                panic!("Expected SRM enable in response");
921            };
922        };
923        expect_response(&mut exec, &mut remote, expectation1, OpCode::Put);
924
925        // Next request sends over some data. Don't expect any response on the remote.
926        let headers2 = HeaderSet::from_header(Header::Body(vec![1, 2, 3, 4, 5]));
927        let put_request2 = RequestPacket::new_put(headers2);
928        send_packet(&mut exec, &mut remote, put_request2);
929        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
930        expect_stream_pending(&mut exec, &mut remote);
931
932        // Next (final) request sends over some data. Expect a response since this is the final
933        // packet.
934        let headers3 = HeaderSet::from_header(Header::EndOfBody(vec![6, 7, 8, 9, 10]));
935        let put_request3 = RequestPacket::new_put_final(headers3);
936        send_packet(&mut exec, &mut remote, put_request3);
937
938        // The entire request is complete and the Obex Server should hand it to the application.
939        // Verify profile received correct data.
940        test_app.set_put_response(Ok(()));
941        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
942        let (rec_data, rec_headers) = test_app.put_data();
943        assert_eq!(rec_data, vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10]);
944        assert!(rec_headers.contains_header(&HeaderIdentifier::Name));
945
946        let expectation = |response: ResponsePacket| {
947            assert_eq!(*response.code(), ResponseCode::Ok);
948            assert!(response.headers().is_empty());
949        };
950        expect_response(&mut exec, &mut remote, expectation, OpCode::PutFinal);
951    }
952
953    #[test_case(Transport::Socket ; "socket")]
954    #[test_case(Transport::Fidl ; "fidl")]
955    #[fuchsia::test]
956    fn delete_request_accepted_by_app_success(transport: Transport) {
957        let mut exec = fasync::TestExecutor::new();
958        let (mut obex_server, test_app, mut remote) =
959            new_obex_server(transport, /*srm=*/ false);
960        obex_server.set_connection_status(ConnectionStatus::connected_no_id());
961        let server_fut = obex_server.run();
962        let mut server_fut = pin!(server_fut);
963        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
964
965        let headers = HeaderSet::from_header(Header::name("foo.txt"));
966        let put_request = RequestPacket::new_put_final(headers);
967        send_packet(&mut exec, &mut remote, put_request);
968
969        // The ObexServer should receive the request and hand it to the profile. Profile accepts.
970        test_app.set_put_response(Ok(()));
971        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
972
973        let expectation = |response: ResponsePacket| {
974            assert_eq!(*response.code(), ResponseCode::Ok);
975            assert!(response.headers().is_empty());
976        };
977        expect_response(&mut exec, &mut remote, expectation, OpCode::PutFinal);
978    }
979
980    #[test_case(Transport::Socket ; "socket")]
981    #[test_case(Transport::Fidl ; "fidl")]
982    #[fuchsia::test]
983    fn receive_packet_exceeds_max_packet_size_rejected(transport: Transport) {
984        let mut exec = fasync::TestExecutor::new();
985        let (mut obex_server, _test_app, mut remote) =
986            new_obex_server(transport, /*srm=*/ false);
987        obex_server.set_connection_status(ConnectionStatus::connected_no_id());
988        // The default max_packet_size before CONNECT is clamped to
989        // MIN_MAX_PACKET_SIZE (255) minimum. Set max_packet_size explicitly to
990        // 255.
991        obex_server.set_max_packet_size(255);
992        let server_fut = obex_server.run();
993        let mut server_fut = pin!(server_fut);
994        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
995
996        let mut oversized_buf = vec![0u8; 300];
997        oversized_buf[0] = (&OpCode::Put).into();
998        oversized_buf[1..3].copy_from_slice(&(300u16).to_be_bytes());
999
1000        let mut fut = remote.send(oversized_buf);
1001        exec.run_until_stalled(&mut fut).expect("write to channel success").expect("write success");
1002        let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
1003
1004        // The ObexServer should reject the oversized packet with BadRequest (0xC0).
1005        let expectation = |response: ResponsePacket| {
1006            assert_eq!(*response.code(), ResponseCode::BadRequest);
1007        };
1008        expect_response(&mut exec, &mut remote, expectation, OpCode::Put);
1009    }
1010
1011    #[test_case(Transport::Socket ; "socket")]
1012    #[test_case(Transport::Fidl ; "fidl")]
1013    #[fuchsia::test]
1014    async fn send_exceeds_max_packet_size_is_error(transport: Transport) {
1015        let (mut obex_server, _test_app, _remote) = new_obex_server(transport, /*srm=*/ false);
1016        obex_server.set_max_packet_size(255);
1017
1018        // Attempt to send a response packet whose encoded length exceeds 255
1019        // bytes.
1020        let large_header = Header::Description("A".repeat(300).into());
1021        let response =
1022            ResponsePacket::new_no_data(ResponseCode::Ok, HeaderSet::from_header(large_header));
1023
1024        let send_result = obex_server.send(response).await;
1025        assert_matches!(send_result, Err(Error::PacketTooLarge { size: 608, max: 255 }));
1026    }
1027}