redevplugin_worker_sdk/
tcp.rs1use crate::api;
2use crate::error::Result;
3use crate::resource::{Handle, MAX_IO_CHUNK_BYTES};
4use serde::{Deserialize, Serialize};
5
6#[derive(Debug, Clone, Serialize)]
7pub struct TcpConnect {
8 pub host: String,
9 pub port: u16,
10 #[serde(default)]
11 #[serde(skip_serializing_if = "Option::is_none")]
12 pub timeout_ms: Option<u32>,
13 #[serde(default)]
14 pub no_delay: bool,
15 #[serde(default)]
16 #[serde(skip_serializing_if = "Option::is_none")]
17 pub keep_alive_ms: Option<u32>,
18}
19
20pub type ConnectOptions = TcpConnect;
21
22#[derive(Deserialize)]
23struct HandleResult {
24 handle: u64,
25}
26
27pub struct TcpStream {
28 handle: Handle,
29}
30
31impl TcpStream {
32 pub fn connect(options: TcpConnect) -> Result<Self> {
33 let opened: HandleResult = api::call("net.tcp.connect", &options)?;
34 Ok(Self {
35 handle: Handle::new(opened.handle)?,
36 })
37 }
38
39 pub fn read(&mut self, capacity: usize) -> Result<(Vec<u8>, u32)> {
40 self.handle.read(capacity)
41 }
42
43 pub fn write_all(&mut self, bytes: &[u8]) -> Result<()> {
44 for chunk in bytes.chunks(MAX_IO_CHUNK_BYTES) {
45 self.handle.write(chunk, 0)?;
46 }
47 Ok(())
48 }
49
50 pub fn shutdown(&mut self, direction: Shutdown) -> Result<()> {
51 #[derive(Serialize)]
52 struct Arguments {
53 handle: u64,
54 direction: Shutdown,
55 }
56 let _: serde_json::Value = api::call(
57 "net.tcp.shutdown",
58 &Arguments {
59 handle: self.handle.id(),
60 direction,
61 },
62 )?;
63 Ok(())
64 }
65
66 pub fn close(mut self) -> Result<()> {
67 self.handle.close()
68 }
69}
70
71#[derive(Debug, Clone, Copy, Serialize)]
72#[serde(rename_all = "snake_case")]
73pub enum Shutdown {
74 Read,
75 Write,
76 Both,
77}
78
79#[derive(Debug, Clone, Serialize)]
80pub struct TcpListen {
81 pub host: String,
82 pub port: u16,
83}
84
85pub type ListenOptions = TcpListen;
86
87#[derive(Deserialize)]
88struct ListenResult {
89 handle: u64,
90 address: String,
91}
92
93pub struct TcpListener {
94 handle: Handle,
95 pub address: String,
96}
97
98impl TcpListener {
99 pub fn listen(options: TcpListen) -> Result<Self> {
100 let opened: ListenResult = api::call("net.tcp.listen", &options)?;
101 Ok(Self {
102 handle: Handle::new(opened.handle)?,
103 address: opened.address,
104 })
105 }
106
107 pub fn accept(&mut self, no_delay: bool, keep_alive_ms: Option<u32>) -> Result<TcpStream> {
108 #[derive(Serialize)]
109 struct Arguments {
110 handle: u64,
111 no_delay: bool,
112 #[serde(skip_serializing_if = "Option::is_none")]
113 keep_alive_ms: Option<u32>,
114 }
115 let opened: HandleResult = api::call(
116 "net.tcp.accept",
117 &Arguments {
118 handle: self.handle.id(),
119 no_delay,
120 keep_alive_ms,
121 },
122 )?;
123 Ok(TcpStream {
124 handle: Handle::new(opened.handle)?,
125 })
126 }
127
128 pub fn close(mut self) -> Result<()> {
129 self.handle.close()
130 }
131}