1use 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
20mod handler;
23pub use handler::{ObexOperationError, ObexServerHandler, new_operation_error};
24
25mod get;
27use get::GetOperation;
28
29mod put;
31use put::PutOperation;
32
33#[derive(Debug)]
35pub enum OperationRequest {
36 SendPackets(Vec<ResponsePacket>),
38 GetApplicationInfo(HeaderSet),
41 GetApplicationData(HeaderSet),
43 PutApplicationData(Vec<u8>, HeaderSet),
45 DeleteApplicationData(HeaderSet),
47 None,
49}
50
51impl OperationRequest {
52 pub fn single_response(packet: ResponsePacket) -> Self {
53 Self::SendPackets(vec![packet])
54 }
55}
56
57#[derive(Debug)]
59pub enum ApplicationResponse {
60 GetInfo(HeaderSet),
63 GetData((Vec<u8>, HeaderSet)),
66 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
87pub trait ServerOperation {
91 fn srm_status(&self) -> SingleResponseMode;
93
94 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 srm_supported_locally && *srm == SingleResponseMode::Enable {
112 Some(SingleResponseMode::Enable)
113 } else {
114 Some(SingleResponseMode::Disable)
116 }
117 }
118
119 fn is_complete(&self) -> bool;
121
122 fn handle_peer_request(&mut self, request: RequestPacket) -> Result<OperationRequest, Error>;
126
127 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 #[default]
143 Initialized,
144 Connected { id: Option<ConnectionIdentifier> },
149 DisconnectReceived,
152}
153
154impl ConnectionStatus {
155 #[cfg(test)]
156 fn connected_no_id() -> Self {
157 Self::Connected { id: None }
158 }
159}
160
161pub struct ObexServer {
165 connected: ConnectionStatus,
167 max_packet_size: u16,
169 active_operation: Option<Box<dyn ServerOperation>>,
176 channel: Channel,
178 type_: TransportType,
180 handler: Box<dyn ObexServerHandler>,
184}
185
186impl ObexServer {
187 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 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 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 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 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 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 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 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 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 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 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 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 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 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, 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, 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 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 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, 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 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 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, 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 test_app.set_response(Err((ResponseCode::Forbidden, HeaderSet::new())));
566 let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
567
568 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, 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 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, 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 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 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, false);
637 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 test_app.set_response(Ok(HeaderSet::new()));
650 let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
651
652 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, false);
668 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 test_app.set_response(Err((ResponseCode::Forbidden, HeaderSet::new())));
680 let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
681
682 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, 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, false);
717 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 let get_request1 =
726 RequestPacket::new_get(HeaderSet::from_header(Header::name("random object")));
727 send_packet(&mut exec, &mut remote, get_request1);
728 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 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 let get_request2 = RequestPacket::new_get_final(HeaderSet::new());
741 send_packet(&mut exec, &mut remote, get_request2);
742
743 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(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, 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 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 let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
780 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, true);
795 obex_server.set_connection_status(ConnectionStatus::connected_no_id());
796 obex_server.set_max_packet_size(20); 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 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 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 let get_request2 = RequestPacket::new_get_final(HeaderSet::new());
826 send_packet(&mut exec, &mut remote, get_request2);
827 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 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 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, false);
863 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 test_app.set_put_response(Ok(()));
879 let _ = exec.run_until_stalled(&mut server_fut).expect_pending("server active");
880 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, 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 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 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 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 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 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, 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 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, false);
987 obex_server.set_connection_status(ConnectionStatus::connected_no_id());
988 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 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, false);
1016 obex_server.set_max_packet_size(255);
1017
1018 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}