1use futures::SinkExt;
6use log::trace;
7
8use crate::client::SrmOperation;
9use crate::error::Error;
10use crate::header::{Header, HeaderIdentifier, HeaderSet, SingleResponseMode};
11use crate::operation::{OpCode, RequestPacket, ResponseCode};
12use crate::transport::ObexTransport;
13
14#[derive(Debug)]
16enum Status {
17 NotStarted(HeaderSet),
21 Started,
23}
24
25#[must_use]
41#[derive(Debug)]
42pub struct PutOperation<'a> {
43 transport: ObexTransport<'a>,
45 status: Status,
47 srm: SingleResponseMode,
50}
51
52impl<'a> PutOperation<'a> {
53 pub fn new(headers: HeaderSet, transport: ObexTransport<'a>) -> Self {
54 let srm = transport.srm_supported().into();
55 Self { transport, status: Status::NotStarted(headers), srm }
56 }
57
58 fn is_started(&self) -> bool {
61 match self.status {
62 Status::NotStarted(_) => false,
63 Status::Started => true,
64 }
65 }
66
67 fn set_started(&mut self) -> Result<(), Error> {
69 match std::mem::replace(&mut self.status, Status::Started) {
70 Status::NotStarted(_) => Ok(()),
71 Status::Started => {
72 Err(Error::other("Attempted to start a PUT operation that was already started"))
73 }
74 }
75 }
76
77 fn combine_with_initial_headers(&mut self, headers: HeaderSet) -> Result<HeaderSet, Error> {
80 let mut initial_headers = match &mut self.status {
81 Status::NotStarted(initial_headers) => std::mem::take(initial_headers),
82 Status::Started => {
83 return Err(Error::other(
84 "Cannot add initial headers when PUT operation already started",
85 ));
86 }
87 };
88 let _ = initial_headers.try_append(headers)?;
89 Ok(initial_headers)
90 }
91
92 fn validate_headers(headers: &HeaderSet) -> Result<(), Error> {
94 if headers.contains_header(&HeaderIdentifier::Body) {
95 return Err(Error::operation(OpCode::Put, "info headers can't contain body"));
96 }
97 if headers.contains_header(&HeaderIdentifier::EndOfBody) {
98 return Err(Error::operation(OpCode::Put, "info headers can't contain end of body"));
99 }
100 Ok(())
101 }
102
103 async fn do_put(&mut self, final_: bool, mut headers: HeaderSet) -> Result<HeaderSet, Error> {
106 let is_started = self.is_started();
107 if !is_started {
108 headers = self.combine_with_initial_headers(headers)?;
109 }
110
111 let srm_active = is_started && self.get_srm() == SingleResponseMode::Enable;
113 let (opcode, request, expected_response_code) = if final_ {
114 (OpCode::PutFinal, RequestPacket::new_put_final(headers), ResponseCode::Ok)
115 } else {
116 (OpCode::Put, RequestPacket::new_put(headers), ResponseCode::Continue)
117 };
118 trace!("Making outgoing PUT request: {request:?}");
119 self.transport.send(request).await?;
120 trace!("Successfully made PUT request");
121 if !is_started {
122 self.set_started()?;
123 }
124 if final_ || !srm_active {
127 let response = self.transport.receive_response(opcode).await?;
128 response.expect_code(opcode, expected_response_code).map(Into::into)
129 } else {
130 Ok(HeaderSet::new())
131 }
132 }
133
134 pub async fn delete(mut self, headers: HeaderSet) -> Result<HeaderSet, Error> {
138 Self::validate_headers(&headers)?;
139 self.do_put(true, headers).await
142 }
143
144 pub async fn write(&mut self, data: &[u8], mut headers: HeaderSet) -> Result<HeaderSet, Error> {
150 Self::validate_headers(&headers)?;
151 let is_first_write = !self.is_started();
152 if is_first_write {
153 self.try_enable_srm(&mut headers)?;
155 }
156 headers.add(Header::Body(data.to_vec()))?;
157 let response_headers = self.do_put(false, headers).await?;
158 if is_first_write {
159 self.check_response_for_srm(&response_headers);
160 }
161 Ok(response_headers)
162 }
163
164 pub async fn write_final(
170 mut self,
171 data: &[u8],
172 mut headers: HeaderSet,
173 ) -> Result<HeaderSet, Error> {
174 Self::validate_headers(&headers)?;
175 headers.add(Header::EndOfBody(data.to_vec()))?;
176 self.do_put(true, headers).await
177 }
178
179 pub async fn terminate(mut self, headers: HeaderSet) -> Result<HeaderSet, Error> {
184 let opcode = OpCode::Abort;
185 if !self.is_started() {
186 return Err(Error::operation(opcode, "can't abort PUT that hasn't started"));
187 }
188 let request = RequestPacket::new_abort(headers);
189 trace!(request:?; "Making outgoing {opcode:?} request");
190 self.transport.send(request).await?;
191 trace!("Successfully made {opcode:?} request");
192 let response = self.transport.receive_response(opcode).await?;
193 response.expect_code(opcode, ResponseCode::Ok).map(Into::into)
194 }
195}
196
197impl SrmOperation for PutOperation<'_> {
198 const OPERATION_TYPE: OpCode = OpCode::Put;
199
200 fn get_srm(&self) -> SingleResponseMode {
201 self.srm
202 }
203
204 fn set_srm(&mut self, mode: SingleResponseMode) {
205 self.srm = mode;
206 }
207}
208
209#[cfg(test)]
210mod tests {
211 use super::*;
212 use bt_channel_test_support::Transport;
213 use test_case::test_case;
214
215 use assert_matches::assert_matches;
216 use async_utils::PollExt;
217 use fuchsia_async as fasync;
218 use std::pin::pin;
219
220 use crate::header::ConnectionIdentifier;
221 use crate::operation::ResponsePacket;
222 use crate::transport::ObexTransportManager;
223 use crate::transport::test_utils::{
224 expect_code, expect_request, expect_request_and_reply, new_manager,
225 };
226
227 fn setup_put_operation(
228 mgr: &ObexTransportManager,
229 initial_headers: Vec<Header>,
230 ) -> PutOperation<'_> {
231 let transport = mgr.try_new_operation().expect("can start operation");
232 PutOperation::new(HeaderSet::from_headers(initial_headers).unwrap(), transport)
233 }
234
235 #[test_case(Transport::Socket ; "socket")]
236 #[test_case(Transport::Fidl ; "fidl")]
237 #[fuchsia::test]
238 fn put_operation_single_chunk_is_ok(transport: Transport) {
239 let mut exec = fasync::TestExecutor::new();
240 let (manager, mut remote) = new_manager(transport, false);
241 let operation =
242 setup_put_operation(&manager, vec![Header::ConnectionId(0x1u32.try_into().unwrap())]);
243
244 let payload = vec![5, 6, 7, 8, 9];
245 let headers =
246 HeaderSet::from_headers(vec![Header::Type("file".into()), Header::name("foobar.txt")])
247 .unwrap();
248 let put_fut = operation.write_final(&payload[..], headers);
249 let mut put_fut = pin!(put_fut);
250 let _ = exec.run_until_stalled(&mut put_fut).expect_pending("waiting for response");
251 let response = ResponsePacket::new_no_data(ResponseCode::Ok, HeaderSet::new());
252 let expectation = |request: RequestPacket| {
253 assert_eq!(*request.code(), OpCode::PutFinal);
254 let headers = HeaderSet::from(request);
255 assert!(headers.contains_header(&HeaderIdentifier::ConnectionId));
256 assert!(!headers.contains_header(&HeaderIdentifier::Body));
257 assert!(headers.contains_headers(&vec![
258 HeaderIdentifier::EndOfBody,
259 HeaderIdentifier::Type,
260 HeaderIdentifier::Name
261 ]));
262 };
263 expect_request_and_reply(&mut exec, &mut remote, expectation, response);
264 let _received_headers = exec
265 .run_until_stalled(&mut put_fut)
266 .expect("response received")
267 .expect("valid response");
268 }
269
270 #[fuchsia::test]
271 fn put_operation_multiple_chunks_is_ok() {
272 let mut exec = fasync::TestExecutor::new();
273 let (manager, mut remote) = new_manager(Transport::Socket, false);
274 let mut operation = setup_put_operation(&manager, vec![]);
275
276 let payload: Vec<u8> = (1..100).collect();
277 for chunk in payload.chunks(20) {
278 let put_fut = operation.write(&chunk[..], HeaderSet::new());
279 let mut put_fut = pin!(put_fut);
280 let _ = exec.run_until_stalled(&mut put_fut).expect_pending("waiting for response");
281 let response = ResponsePacket::new_no_data(ResponseCode::Continue, HeaderSet::new());
282 let expectation = |request: RequestPacket| {
283 assert_eq!(*request.code(), OpCode::Put);
284 let headers = HeaderSet::from(request);
285 assert!(headers.contains_header(&HeaderIdentifier::Body));
286 };
287 expect_request_and_reply(&mut exec, &mut remote, expectation, response);
288 let _received_headers = exec
289 .run_until_stalled(&mut put_fut)
290 .expect("response received")
291 .expect("valid response");
292 }
293
294 let put_final_fut = operation.write_final(&[], HeaderSet::new());
296 let mut put_final_fut = pin!(put_final_fut);
297 let _ = exec.run_until_stalled(&mut put_final_fut).expect_pending("waiting for response");
298 let response = ResponsePacket::new_no_data(ResponseCode::Ok, HeaderSet::new());
299 let expectation = |request: RequestPacket| {
300 assert_eq!(*request.code(), OpCode::PutFinal);
301 let headers = HeaderSet::from(request);
302 assert!(headers.contains_header(&HeaderIdentifier::EndOfBody));
303 };
304 expect_request_and_reply(&mut exec, &mut remote, expectation, response);
305 let _ = exec
306 .run_until_stalled(&mut put_final_fut)
307 .expect("response received")
308 .expect("valid response");
309 }
310
311 #[fuchsia::test]
312 fn put_operation_delete_is_ok() {
313 let mut exec = fasync::TestExecutor::new();
314 let (manager, mut remote) = new_manager(Transport::Socket, false);
315 let operation = setup_put_operation(&manager, vec![]);
316
317 let headers = HeaderSet::from_headers(vec![
318 Header::Description("deleting file".into()),
319 Header::name("foobar.txt"),
320 ])
321 .unwrap();
322 let put_fut = operation.delete(headers);
323 let mut put_fut = pin!(put_fut);
324 let _ = exec.run_until_stalled(&mut put_fut).expect_pending("waiting for response");
325 let response = ResponsePacket::new_no_data(ResponseCode::Ok, HeaderSet::new());
326 let expectation = |request: RequestPacket| {
327 assert_eq!(*request.code(), OpCode::PutFinal);
328 let headers = HeaderSet::from(request);
329 assert!(!headers.contains_header(&HeaderIdentifier::Body));
330 assert!(!headers.contains_header(&HeaderIdentifier::EndOfBody));
331 };
332 expect_request_and_reply(&mut exec, &mut remote, expectation, response);
333 let _ = exec
334 .run_until_stalled(&mut put_fut)
335 .expect("response received")
336 .expect("valid response");
337 }
338
339 #[fuchsia::test]
340 fn put_operation_terminate_success() {
341 let mut exec = fasync::TestExecutor::new();
342 let (manager, mut remote) = new_manager(Transport::Socket, false);
343 let mut operation = setup_put_operation(&manager, vec![]);
344
345 {
347 let put_fut = operation.write(&[1, 2, 3, 4, 5], HeaderSet::new());
348 let mut put_fut = pin!(put_fut);
349 let _ = exec.run_until_stalled(&mut put_fut).expect_pending("waiting for response");
350 let response = ResponsePacket::new_no_data(ResponseCode::Continue, HeaderSet::new());
351 expect_request_and_reply(&mut exec, &mut remote, expect_code(OpCode::Put), response);
352 let _received_headers = exec
353 .run_until_stalled(&mut put_fut)
354 .expect("response received")
355 .expect("valid response");
356 }
357
358 let terminate_fut = operation.terminate(HeaderSet::new());
360 let mut terminate_fut = pin!(terminate_fut);
361 let _ = exec.run_until_stalled(&mut terminate_fut).expect_pending("waiting for response");
362 let response = ResponsePacket::new_no_data(ResponseCode::Ok, HeaderSet::new());
363 expect_request_and_reply(&mut exec, &mut remote, expect_code(OpCode::Abort), response);
364 let _received_headers = exec
365 .run_until_stalled(&mut terminate_fut)
366 .expect("response received")
367 .expect("valid response");
368 }
369
370 #[fuchsia::test]
371 async fn put_with_body_header_is_error() {
372 let (manager, _remote) = new_manager(Transport::Socket, false);
373 let mut operation = setup_put_operation(&manager, vec![]);
374
375 let payload = vec![1, 2, 3];
376 let body_headers = HeaderSet::from_headers(vec![
379 Header::Body(payload.clone()),
380 Header::name("foobar.txt"),
381 ])
382 .unwrap();
383 let result = operation.write(&payload[..], body_headers.clone()).await;
384 assert_matches!(result, Err(Error::OperationError { .. }));
385
386 let eob_headers = HeaderSet::from_headers(vec![
388 Header::EndOfBody(payload.clone()),
389 Header::name("foobar1.txt"),
390 ])
391 .unwrap();
392 let result = operation.write(&payload[..], eob_headers.clone()).await;
393 assert_matches!(result, Err(Error::OperationError { .. }));
394 }
395
396 #[fuchsia::test]
397 async fn delete_with_body_header_is_error() {
398 let (manager, _remote) = new_manager(Transport::Socket, false);
399
400 let payload = vec![1, 2, 3];
401 let operation = setup_put_operation(&manager, vec![]);
403 let body_headers = HeaderSet::from_headers(vec![
404 Header::Body(payload.clone()),
405 Header::name("foobar.txt"),
406 ])
407 .unwrap();
408 let result = operation.delete(body_headers).await;
409 assert_matches!(result, Err(Error::OperationError { .. }));
410
411 let operation = setup_put_operation(&manager, vec![]);
413 let eob_headers = HeaderSet::from_headers(vec![
414 Header::EndOfBody(payload.clone()),
415 Header::name("foobar1.txt"),
416 ])
417 .unwrap();
418 let result = operation.delete(eob_headers).await;
419 assert_matches!(result, Err(Error::OperationError { .. }));
420 }
421
422 #[fuchsia::test]
423 async fn put_operation_terminate_before_start_error() {
424 let (manager, _remote) = new_manager(Transport::Socket, false);
425 let operation = setup_put_operation(&manager, vec![]);
426
427 let headers = HeaderSet::from_header(Header::Description("terminating test".into()));
429 let terminate_result = operation.terminate(headers).await;
430 assert_matches!(terminate_result, Err(Error::OperationError { .. }));
431 }
432
433 #[fuchsia::test]
434 fn put_operation_srm_enabled_is_ok() {
435 let mut exec = fasync::TestExecutor::new();
436 let (manager, mut remote) = new_manager(Transport::Socket, true);
437 let mut operation = setup_put_operation(&manager, vec![]);
438
439 {
440 let first_buf = [1, 2, 3];
441 let put_fut = operation.write(&first_buf[..], HeaderSet::new());
443 let mut put_fut = pin!(put_fut);
444 let _ = exec.run_until_stalled(&mut put_fut).expect_pending("waiting for response");
445
446 let expectation = |request: RequestPacket| {
449 assert_eq!(*request.code(), OpCode::Put);
450 let headers = HeaderSet::from(request);
451 assert!(headers.contains_header(&HeaderIdentifier::Body));
452 assert!(headers.contains_header(&HeaderIdentifier::SingleResponseMode));
453 };
454 let response_headers = HeaderSet::from_header(SingleResponseMode::Enable.into());
455 let response = ResponsePacket::new_no_data(ResponseCode::Continue, response_headers);
456 expect_request_and_reply(&mut exec, &mut remote, expectation, response);
457 let _received_headers = exec
458 .run_until_stalled(&mut put_fut)
459 .expect("response received")
460 .expect("valid response");
461 }
462 assert_eq!(operation.srm, SingleResponseMode::Enable);
464 {
466 let second_buf = [4, 5, 6];
467 let put_fut2 = operation.write(&second_buf[..], HeaderSet::new());
468 let mut put_fut2 = pin!(put_fut2);
469 let _ = exec
470 .run_until_stalled(&mut put_fut2)
471 .expect("ready without peer response")
472 .expect("success");
473 let expectation = |request: RequestPacket| {
474 assert_eq!(*request.code(), OpCode::Put);
475 let headers = HeaderSet::from(request);
476 assert!(headers.contains_header(&HeaderIdentifier::Body));
477 assert!(!headers.contains_header(&HeaderIdentifier::SingleResponseMode));
478 };
479 expect_request(&mut exec, &mut remote, expectation);
480 }
481
482 let put_final_fut = operation.write_final(&[], HeaderSet::new());
484 let mut put_final_fut = pin!(put_final_fut);
485 let _ = exec.run_until_stalled(&mut put_final_fut).expect_pending("waiting for response");
486 let response = ResponsePacket::new_no_data(ResponseCode::Ok, HeaderSet::new());
487 let expectation = |request: RequestPacket| {
488 assert_eq!(*request.code(), OpCode::PutFinal);
489 let headers = HeaderSet::from(request);
490 assert!(headers.contains_header(&HeaderIdentifier::EndOfBody));
491 };
492 expect_request_and_reply(&mut exec, &mut remote, expectation, response);
493 let _ = exec
494 .run_until_stalled(&mut put_final_fut)
495 .expect("response received")
496 .expect("valid response");
497 }
498
499 #[fuchsia::test]
500 fn client_disable_srm_mid_operation_is_ignored() {
501 let mut exec = fasync::TestExecutor::new();
502 let (manager, mut remote) = new_manager(Transport::Socket, true);
503 let mut operation = setup_put_operation(&manager, vec![]);
504 if let Status::NotStarted(_) = &mut operation.status {
506 let _ = operation.set_started().unwrap();
507 } else {
508 panic!("At this point operation not started");
509 };
510 assert_eq!(operation.srm, SingleResponseMode::Enable);
512
513 {
515 let headers = HeaderSet::from_header(SingleResponseMode::Disable.into());
516 let put_fut = operation.write(&[], headers);
517 let mut put_fut = pin!(put_fut);
518 let _ = exec
519 .run_until_stalled(&mut put_fut)
520 .expect("ready without peer response")
521 .expect("success");
522 let expectation = |request: RequestPacket| {
523 assert_eq!(*request.code(), OpCode::Put);
524 };
525 expect_request(&mut exec, &mut remote, expectation);
526 }
527 assert_eq!(operation.srm, SingleResponseMode::Enable);
529 }
530
531 #[fuchsia::test]
532 fn application_select_srm_success() {
533 let _exec = fasync::TestExecutor::new();
534 let (manager, _remote) = new_manager(Transport::Socket, false);
535 let mut operation = setup_put_operation(&manager, vec![]);
536 assert_eq!(operation.srm, SingleResponseMode::Disable);
537 let mut headers = HeaderSet::from_header(SingleResponseMode::Disable.into());
539 assert_matches!(operation.try_enable_srm(&mut headers), Ok(()));
540 assert_eq!(operation.srm, SingleResponseMode::Disable);
541
542 let (manager, _remote) = new_manager(Transport::Socket, true);
544 let mut operation = setup_put_operation(&manager, vec![]);
545 assert_eq!(operation.srm, SingleResponseMode::Enable);
546 let mut headers = HeaderSet::from_header(SingleResponseMode::Disable.into());
547 assert_matches!(operation.try_enable_srm(&mut headers), Ok(()));
548 assert_eq!(operation.srm, SingleResponseMode::Disable);
549
550 let (manager, _remote) = new_manager(Transport::Socket, true);
552 let mut operation = setup_put_operation(&manager, vec![]);
553 assert_eq!(operation.srm, SingleResponseMode::Enable);
554 let mut headers = HeaderSet::from_header(SingleResponseMode::Enable.into());
555 assert_matches!(operation.try_enable_srm(&mut headers), Ok(()));
556 assert_eq!(operation.srm, SingleResponseMode::Enable);
557 }
558
559 #[fuchsia::test]
560 fn application_enable_srm_when_not_supported_is_error() {
561 let _exec = fasync::TestExecutor::new();
562 let (manager, _remote) = new_manager(Transport::Socket, false);
563 let mut operation = setup_put_operation(&manager, vec![]);
564 assert_eq!(operation.srm, SingleResponseMode::Disable);
565 let mut headers = HeaderSet::from_header(SingleResponseMode::Enable.into());
566 assert_matches!(operation.try_enable_srm(&mut headers), Err(Error::SrmNotSupported));
567 assert_eq!(operation.srm, SingleResponseMode::Disable);
568 }
569
570 #[fuchsia::test]
571 fn peer_srm_response() {
572 let _exec = fasync::TestExecutor::new();
573 let (manager, _remote) = new_manager(Transport::Socket, false);
574 let mut operation = setup_put_operation(&manager, vec![]);
575 let headers = HeaderSet::from_header(SingleResponseMode::Enable.into());
577 operation.check_response_for_srm(&headers);
578 assert_eq!(operation.srm, SingleResponseMode::Disable);
579 let headers = HeaderSet::from_header(SingleResponseMode::Disable.into());
581 operation.check_response_for_srm(&headers);
582 assert_eq!(operation.srm, SingleResponseMode::Disable);
583
584 let (manager, _remote) = new_manager(Transport::Socket, true);
585 let mut operation = setup_put_operation(&manager, vec![]);
586 let headers = HeaderSet::from_header(SingleResponseMode::Enable.into());
588 operation.check_response_for_srm(&headers);
589 assert_eq!(operation.srm, SingleResponseMode::Enable);
590 let headers = HeaderSet::from_header(SingleResponseMode::Disable.into());
592 operation.check_response_for_srm(&headers);
593 assert_eq!(operation.srm, SingleResponseMode::Disable);
594
595 let (manager, _remote) = new_manager(Transport::Socket, true);
596 let mut operation = setup_put_operation(&manager, vec![]);
597 operation.check_response_for_srm(&HeaderSet::new());
599 assert_eq!(operation.srm, SingleResponseMode::Disable);
600 }
601
602 #[fuchsia::test]
603 fn put_with_connection_id_already_set_is_error() {
604 let mut exec = fasync::TestExecutor::new();
605 let (manager, _remote) = new_manager(Transport::Socket, false);
606 let mut operation =
608 setup_put_operation(&manager, vec![Header::ConnectionId(ConnectionIdentifier(5))]);
609
610 let write_headers = HeaderSet::from_header(Header::ConnectionId(ConnectionIdentifier(10)));
611 let write_fut = operation.write(&[1, 2, 3], write_headers);
612 let mut write_fut = pin!(write_fut);
613 let result = exec.run_until_stalled(&mut write_fut).expect("finished with error");
614 assert_matches!(result, Err(Error::AlreadyExists(_)));
615 }
616}