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
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
use std::net::{IpAddr, SocketAddr};
use std::pin::Pin;
use std::task::{Context, Poll};

use anyhow::{Error, Result};
use pin_project::pin_project;
use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, ReadBuf};
use tokio::net::{TcpListener, TcpStream};

use crate::{WaitOptions, Waitable};

/// Listens on a specific IP Address and Port using TCP protocol
pub struct TcpWaiter {
    pub addr: IpAddr,
    pub port: u16,
}

impl TcpWaiter {
    pub fn new(addr: IpAddr, port: u16) -> Self {
        Self { addr, port }
    }

    pub fn socket(&self) -> SocketAddr {
        SocketAddr::new(self.addr, self.port)
    }
}

impl Waitable for TcpWaiter {
    async fn wait(self, _: WaitOptions) -> Result<()> {
        let tcp_listener = TcpListener::bind(self.socket()).await?;
        let (socket, _) = tcp_listener.accept().await?;
        let mut socket = PacketExtractor::<8>::read(socket).await?;

        tokio::spawn(async move {
            let mut buf = vec![0; 1024];

            loop {
                let n = socket
                    .read(&mut buf)
                    .await
                    .expect("failed to read data from socket");

                if n == 0 {
                    // socket closed
                    return;
                }
            }
        })
        .await
        .map_err(|err| Error::msg(err.to_string()))?;

        Ok(())
    }
}

#[pin_project]
pub struct PacketExtractor<const B: usize> {
    pub header: [u8; B],
    pub forwarded: usize,
    #[pin]
    pub socket: TcpStream,
}

impl<const B: usize> PacketExtractor<B> {
    pub async fn read(socket: TcpStream) -> Result<Self> {
        let mut extractor = Self {
            header: [0; B],
            forwarded: 0,
            socket,
        };

        extractor.socket.read_exact(&mut extractor.header).await?;

        Ok(extractor)
    }

    pub fn get_header(&mut self) -> &[u8; B] {
        &self.header
    }
}

impl<const B: usize> AsyncRead for PacketExtractor<B> {
    fn poll_read(
        self: Pin<&mut Self>,
        cx: &mut Context<'_>,
        buff: &mut ReadBuf<'_>,
    ) -> Poll<std::io::Result<()>> {
        let extractor = self.project();

        if *extractor.forwarded < extractor.header.len() {
            let leftover = &extractor.header[*extractor.forwarded..];
            let num_forward_now = leftover.len().min(buff.remaining());
            let forward = &leftover[..num_forward_now];

            buff.put_slice(forward);
            *extractor.forwarded += num_forward_now;

            return Poll::Ready(Ok(()));
        }

        extractor.socket.poll_read(cx, buff)
    }
}

impl<const B: usize> AsyncWrite for PacketExtractor<B> {
    fn poll_write(
        self: Pin<&mut Self>,
        cx: &mut Context<'_>,
        buff: &[u8],
    ) -> Poll<Result<usize, std::io::Error>> {
        let extractor = self.project();
        extractor.socket.poll_write(cx, buff)
    }

    fn poll_flush(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<(), std::io::Error>> {
        let extractor = self.project();
        extractor.socket.poll_flush(cx)
    }

    fn poll_shutdown(
        self: Pin<&mut Self>,
        cx: &mut Context<'_>,
    ) -> Poll<Result<(), std::io::Error>> {
        let extractor = self.project();
        extractor.socket.poll_shutdown(cx)
    }
}