Skip to main content

imsg_map/client/
mod.rs

1//! MAP client — OBEX session setup, SETPATH sequencing, and request dispatch.
2//!
3//! Split across `nav` (folder navigation), `messages` (message CRUD), and `control`
4//! (folder listing, notifications, session teardown); session setup and the shared
5//! `recv`/`collect_body` primitives live here.
6
7mod control;
8mod messages;
9mod nav;
10
11use bytes::Bytes;
12use futures::{SinkExt, StreamExt};
13use obex_core::client::ObexClient;
14use obex_core::{wrap, ObexTransport};
15use tokio::io::{AsyncRead, AsyncWrite};
16
17use crate::MapError;
18
19const MAP_UUID: [u8; 16] = [
20    0xbb, 0x58, 0x2b, 0x40, 0x42, 0x0c, 0x11, 0xdb, 0xb0, 0xde, 0x08, 0x00, 0x20, 0x0c, 0x9a, 0x66,
21];
22
23const MAX_BODY_BYTES: usize = 4 * 1024 * 1024;
24
25/// MAP session over an OBEX transport. Owns the OBEX state machine and the framed I/O.
26///
27/// Obtain via [`MapClient::connect`].
28pub struct MapClient<T> {
29    obex: ObexClient,
30    transport: ObexTransport<T>,
31    // SETPATH depth confirmed by server; 0 = root, max 3 (telecom/msg/<folder>).
32    depth: u8,
33}
34
35impl<T: AsyncRead + AsyncWrite + Unpin> MapClient<T> {
36    /// Sends the MAP UUID as the OBEX `Target` header and validates the server's response.
37    ///
38    /// # Errors
39    ///
40    /// Returns [`MapError`] if the transport fails, packet encoding fails, the server rejects
41    /// the connection, or the response omits the `ConnectionId` header.
42    pub async fn connect(stream: T) -> Result<Self, MapError> {
43        let mut transport = wrap(stream);
44        let mut obex = ObexClient::new();
45        let req = ObexClient::connect_request(&MAP_UUID, None)?;
46        transport.send(req).await?;
47        let rsp = Self::recv(&mut transport).await?;
48        obex.handle_connect_response(&rsp)?;
49        Ok(Self { obex, transport, depth: 0 })
50    }
51
52    async fn recv(transport: &mut ObexTransport<T>) -> Result<Bytes, MapError> {
53        transport.next().await.ok_or(MapError::UnexpectedEof)?.map_err(MapError::Transport)
54    }
55
56    async fn collect_body(&mut self) -> Result<Vec<u8>, MapError> {
57        let mut body = Vec::with_capacity(512);
58        loop {
59            let rsp_bytes = Self::recv(&mut self.transport).await?;
60            let rsp = ObexClient::parse_response(&rsp_bytes)?;
61            if rsp.opcode.is_continue() {
62                if let Some(chunk) = rsp.body_payload() {
63                    let new_len =
64                        body.len().checked_add(chunk.len()).ok_or(MapError::ResponseTooLarge)?;
65                    if new_len > MAX_BODY_BYTES {
66                        return Err(MapError::ResponseTooLarge);
67                    }
68                    body.extend_from_slice(chunk);
69                }
70                let cont = self.obex.get_continue_request()?;
71                self.transport.send(cont).await?;
72            } else if rsp.opcode.is_ok() {
73                if let Some(chunk) = rsp.body_payload() {
74                    let new_len =
75                        body.len().checked_add(chunk.len()).ok_or(MapError::ResponseTooLarge)?;
76                    if new_len > MAX_BODY_BYTES {
77                        return Err(MapError::ResponseTooLarge);
78                    }
79                    body.extend_from_slice(chunk);
80                }
81                break;
82            } else {
83                return Err(MapError::ServerError(rsp.opcode.to_byte()));
84            }
85        }
86        Ok(body)
87    }
88}