1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
use super::*;

pub use crate::public::depth::LimitOrder;
use crate::public::depth::RawLimitOrder;

#[serde_as]
#[derive(Deserialize, Debug)]
struct RawResponse {
    a: Vec<RawLimitOrder>,
    b: Vec<RawLimitOrder>,
    #[serde_as(as = "TimestampMilliSeconds")]
    t: NaiveDateTime,
    #[serde_as(as = "DisplayFromStr")]
    s: u64,
}

#[derive(Debug)]
pub struct DepthDiff {
    pub asks: Vec<LimitOrder>,
    pub bids: Vec<LimitOrder>,
    pub timestamp: NaiveDateTime,
    pub sequence_id: u64,
}

#[derive(TypedBuilder)]
pub struct Params {
    pair: Pair,
}

pub async fn connect(
    params: Params,
) -> anyhow::Result<impl tokio_stream::Stream<Item = DepthDiff>> {
    use tokio_stream::StreamExt;

    let pair = params.pair;
    let room_id = format!("depth_diff_{pair}");
    let raw = do_connect::<RawResponse>(&room_id).await?;
    let st = raw.map(|x| DepthDiff {
        asks: x.a.into_iter().map(LimitOrder::new).collect(),
        bids: x.b.into_iter().map(LimitOrder::new).collect(),
        timestamp: x.t,
        sequence_id: x.s,
    });
    Ok(st)
}

#[cfg(test)]
mod tests {
    use super::*;

    #[tokio::test]
    async fn test() -> anyhow::Result<()> {
        use futures_util::{pin_mut, StreamExt};

        let params = Params::builder().pair(Pair(XRP, JPY)).build();
        let st = connect(params).await?;
        pin_mut!(st);
        while let Some(x) = st.next().await {
            dbg!(&x);
            break;
        }
        Ok(())
    }
}