Skip to main content

bt_obex/client/
put.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 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/// Represents the status of the PUT operation.
15#[derive(Debug)]
16enum Status {
17    /// First write call has not been completed yet.
18    /// Holds the initial headers that need to be included in the
19    /// first write call.
20    NotStarted(HeaderSet),
21    /// First write has been completed and the operation is ongoing.
22    Started,
23}
24
25/// Represents an in-progress PUT Operation.
26/// Defined in OBEX 1.5 Section 3.4.3.
27///
28/// Example Usage:
29/// ```
30/// let obex_client = ObexClient::new(..);
31/// let put_operation = obex_client.put()?;
32/// let user_data: Vec<u8> = vec![];
33/// for user_data_chunk in user_data.chunks(50) {
34///   let received_headers = put_operation.write(&user_data_chunk[..], HeaderSet::new()).await?;
35/// }
36/// // `PutOperation::write_final` must be called before it is dropped. An empty payload is OK.
37/// let final_headers = put_operation.write_final(&[], HeaderSet::new()).await?;
38/// // PUT operation is complete and `put_operation` is consumed.
39/// ```
40#[must_use]
41#[derive(Debug)]
42pub struct PutOperation<'a> {
43    /// The L2CAP or RFCOMM connection to the remote peer.
44    transport: ObexTransport<'a>,
45    /// Status of the operation.
46    status: Status,
47    /// The status of SRM for this operation. By default, SRM will be enabled if the transport
48    /// supports it. However, it may be disabled if the peer requests to disable it.
49    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    /// Returns true by checking whether the initial headers were taken
59    /// out for the first put operation.
60    fn is_started(&self) -> bool {
61        match self.status {
62            Status::NotStarted(_) => false,
63            Status::Started => true,
64        }
65    }
66
67    /// Sets the operation as started.
68    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    /// Returns the HeaderSet that takes the initial headers and
78    /// combines them with the input headers.
79    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    /// Returns Error if the `headers` contain non-informational OBEX Headers.
93    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    /// Attempts to initiate a PUT operation with the `final_` bit set.
104    /// Returns the peer response headers on success, Error otherwise.
105    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        // SRM is considered active if this is a subsequent PUT request & the transport supports it.
112        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        // Expect a response if this is the final PUT request or if SRM is inactive, in which case
125        // every request must be responded to.
126        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    /// Attempts to delete an object from the remote OBEX server specified by the provided
135    /// `headers`.
136    /// Returns the informational headers from the peer response on success, Error otherwise.
137    pub async fn delete(mut self, headers: HeaderSet) -> Result<HeaderSet, Error> {
138        Self::validate_headers(&headers)?;
139        // No Body or EndOfBody Headers are included in a delete request.
140        // See OBEX 1.5 Section 3.4.3.6.
141        self.do_put(true, headers).await
142    }
143
144    /// Attempts to write the `data` object to the remote OBEX server.
145    /// Returns the informational headers from the peer response on success, Error otherwise.
146    /// The returned informational headers will be empty if Single Response Mode is enabled for the
147    /// operation. Only the final write request (`Self::write_final`) will potentially return a
148    /// non-empty set of headers.
149    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            // Try to enable SRM if this is the first packet of the operation.
154            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    /// Attempts to write the final `data` object to the remote OBEX server.
165    /// This must be called before the PutOperation object is dropped.
166    /// Returns the informational headers from the peer response on success, Error otherwise.
167    ///
168    /// The PUT operation is considered complete after this.
169    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    /// Request to terminate a multi-packet PUT request early.
180    /// Returns the informational headers from the peer response on success, Error otherwise.
181    /// If Error is returned, there are no guarantees about the synchronization between the local
182    /// OBEX client and remote OBEX server.
183    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, /* srm_supported */ 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, /* srm_supported */ 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        // Can send final response that is empty to complete the operation.
295        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, /* srm_supported */ 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, /* srm_supported */ false);
343        let mut operation = setup_put_operation(&manager, vec![]);
344
345        // Write the first chunk of data to "start" the operation.
346        {
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        // Terminating early should be Ok - peer acknowledges.
359        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, /* srm_supported */ false);
373        let mut operation = setup_put_operation(&manager, vec![]);
374
375        let payload = vec![1, 2, 3];
376        // The payload should only be included as an argument. All other headers must be
377        // informational.
378        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        // EndOfBody header is also an Error.
387        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, /* srm_supported */ false);
399
400        let payload = vec![1, 2, 3];
401        // Body shouldn't be included in delete.
402        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        // EndOfBody shouldn't be included in delete.
412        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, /* srm_supported */ false);
425        let operation = setup_put_operation(&manager, vec![]);
426
427        // Trying to terminate early doesn't work as the operation has not started.
428        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, /* srm_supported */ true);
437        let mut operation = setup_put_operation(&manager, vec![]);
438
439        {
440            let first_buf = [1, 2, 3];
441            // Even though the input headers are empty, we should prefer to enable SRM.
442            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            // Expect the outgoing request with the SRM Header. Peer responds positively with a SRM
447            // Enable response.
448            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        // At this point SRM is enabled for the duration of the operation.
463        assert_eq!(operation.srm, SingleResponseMode::Enable);
464        // Second write doesn't require a response.
465        {
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        // Only the final write request will result in a response.
483        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, /* srm_supported */ true);
503        let mut operation = setup_put_operation(&manager, vec![]);
504        // Pretend first write happened already by manually setting the operation as started.
505        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        // SRM is enabled for the duration of the operation.
511        assert_eq!(operation.srm, SingleResponseMode::Enable);
512
513        // Client tries to disable SRM in a subsequent write attempt. Ignored.
514        {
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        // SRM is still enabled.
528        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, /* srm_supported */ false);
535        let mut operation = setup_put_operation(&manager, vec![]);
536        assert_eq!(operation.srm, SingleResponseMode::Disable);
537        // The application requesting to disable SRM when it isn't supported is OK.
538        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        // The application requesting to disable SRM when it is supported is OK.
543        let (manager, _remote) = new_manager(Transport::Socket, /* srm_supported */ 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        // The application requesting to enable SRM when it is supported is OK.
551        let (manager, _remote) = new_manager(Transport::Socket, /* srm_supported */ 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, /* srm_supported */ 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, /* srm_supported */ false);
574        let mut operation = setup_put_operation(&manager, vec![]);
575        // An enable response from the peer when SRM is disabled locally should not enable SRM.
576        let headers = HeaderSet::from_header(SingleResponseMode::Enable.into());
577        operation.check_response_for_srm(&headers);
578        assert_eq!(operation.srm, SingleResponseMode::Disable);
579        // A disable response from the peer when SRM is disabled locally is a no-op.
580        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, /* srm_supported */ true);
585        let mut operation = setup_put_operation(&manager, vec![]);
586        // An enable response from the peer when SRM is enable is a no-op.
587        let headers = HeaderSet::from_header(SingleResponseMode::Enable.into());
588        operation.check_response_for_srm(&headers);
589        assert_eq!(operation.srm, SingleResponseMode::Enable);
590        // A disable response from the peer when SRM is enabled should disable SRM.
591        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, /* srm_supported */ true);
596        let mut operation = setup_put_operation(&manager, vec![]);
597        // A response with no SRM header should be treated like a disable request.
598        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, /* srm_supported */ false);
606        // The initial operation contains a ConnectionId header which was negotiated during CONNECT.
607        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}