mail4agent_messenger_shell/
send.rs1use serde::{Deserialize, Serialize};
15use std::ffi::OsStr;
16use std::io::{BufRead, BufReader, Read, Write};
17use std::path::{Path, PathBuf};
18use std::time::Duration;
19
20use crate::ipc::SendStream;
21
22pub const SEND_SOCK_ENV: &str = "M4A_SEND_SOCK";
24
25pub const DEFAULT_SOCK_NAME: &str = "web-client.sock";
27
28pub const ENV_FILE_ENV: &str = "M4A_ENV_FILE";
34
35pub const MAX_SEND_BYTES: usize = 16 * 1024;
37
38#[derive(Debug, Clone, Serialize, Deserialize)]
40pub struct SendRequest {
41 #[serde(rename = "as")]
43 pub as_nick: String,
44 pub to: String,
46 pub text: String,
48}
49
50#[derive(Debug, Clone, Default, Serialize, Deserialize)]
52pub struct SendReply {
53 pub ok: bool,
55 #[serde(default, skip_serializing_if = "Option::is_none")]
57 pub room: Option<String>,
58 #[serde(default, skip_serializing_if = "Option::is_none")]
60 pub event_id: Option<String>,
61 #[serde(default, skip_serializing_if = "Option::is_none")]
63 pub error: Option<String>,
64}
65
66impl SendReply {
67 pub fn failed(error: impl Into<String>) -> Self {
69 Self {
70 error: Some(error.into()),
71 ..Self::default()
72 }
73 }
74}
75
76pub fn send_sock_path(get: impl FnMut(&str) -> Option<String>, store_root: &Path) -> PathBuf {
79 send_sock_path_named(get, store_root, DEFAULT_SOCK_NAME)
80}
81
82pub fn send_sock_path_named(
85 mut get: impl FnMut(&str) -> Option<String>,
86 store_root: &Path,
87 default_name: &str,
88) -> PathBuf {
89 get(SEND_SOCK_ENV)
90 .map(PathBuf::from)
91 .unwrap_or_else(|| store_root.join(default_name))
92}
93
94pub fn load_env_file() -> Option<PathBuf> {
100 load_env_file_named("web-client.env")
101}
102
103pub fn load_env_file_named(default_name: &str) -> Option<PathBuf> {
108 let path = std::env::var(ENV_FILE_ENV)
109 .ok()
110 .filter(|value| !value.is_empty())
111 .map(PathBuf::from)
112 .or_else(|| {
113 home_dir_from(
114 std::env::var_os("HOME").as_deref(),
115 std::env::var_os("USERPROFILE").as_deref(),
116 )
117 .map(|home| home.join(".config/mail4agent").join(default_name))
118 })?;
119 let text = std::fs::read_to_string(&path).ok()?;
120 for (key, value) in parse_env_lines(&text) {
121 if std::env::var_os(&key).is_none() {
122 std::env::set_var(key, value);
123 }
124 }
125 Some(path)
126}
127
128fn home_dir_from(home: Option<&OsStr>, userprofile: Option<&OsStr>) -> Option<PathBuf> {
132 let chosen = home
133 .filter(|value| !value.is_empty())
134 .or_else(|| userprofile.filter(|value| !value.is_empty()))?;
135 Some(PathBuf::from(chosen))
136}
137
138fn parse_env_lines(text: &str) -> Vec<(String, String)> {
139 text.lines()
140 .filter_map(|line| {
141 let line = line.trim();
142 if line.is_empty() || line.starts_with('#') {
143 return None;
144 }
145 let line = line.strip_prefix("export ").unwrap_or(line);
146 let (key, value) = line.split_once('=')?;
147 let key = key.trim();
148 if key.is_empty() || !key.bytes().all(|b| b.is_ascii_alphanumeric() || b == b'_') {
149 return None;
150 }
151 let value = value.trim();
152 let value = value
153 .strip_prefix('"')
154 .and_then(|v| v.strip_suffix('"'))
155 .or_else(|| value.strip_prefix('\'').and_then(|v| v.strip_suffix('\'')))
156 .unwrap_or(value);
157 Some((key.to_string(), value.to_string()))
158 })
159 .collect()
160}
161
162pub fn send_via_socket(
167 sock: &Path,
168 request: &SendRequest,
169 wait: Duration,
170) -> std::io::Result<SendReply> {
171 let mut stream = SendStream::connect(sock)?;
172 stream.set_read_timeout(Some(wait))?;
173 stream.set_write_timeout(Some(Duration::from_secs(5)))?;
174 let mut line = serde_json::to_vec(request)?;
175 line.push(b'\n');
176 stream.write_all(&line)?;
177 stream.flush()?;
178 let mut reader = BufReader::new(stream);
179 let mut answer = String::new();
180 reader.read_line(&mut answer)?;
181 if answer.trim().is_empty() {
182 return Err(std::io::Error::new(
183 std::io::ErrorKind::UnexpectedEof,
184 "client closed the socket without an answer",
185 ));
186 }
187 serde_json::from_str(answer.trim())
188 .map_err(|err| std::io::Error::new(std::io::ErrorKind::InvalidData, err))
189}
190
191pub(crate) enum Incoming {
193 Send(SendRequest),
194 Cmd(crate::CmdRequest),
195}
196
197pub(crate) fn read_incoming(stream: &mut SendStream) -> Result<Incoming, String> {
199 stream.set_nonblocking(false).map_err(|err| err.to_string())?;
200 stream.set_read_timeout(Some(Duration::from_secs(3))).map_err(|err| err.to_string())?;
201 let mut buf = Vec::new();
202 let mut limited = stream.take((MAX_SEND_BYTES * 2 + 1024) as u64);
203 let mut reader = BufReader::new(&mut limited);
204 reader.read_until(b'\n', &mut buf).map_err(|err| err.to_string())?;
205 let value: serde_json::Value =
206 serde_json::from_slice(&buf).map_err(|_| "request is not a JSON line".to_string())?;
207 if value.get("cmd").is_some() {
208 let cmd: crate::CmdRequest =
209 serde_json::from_value(value).map_err(|err| format!("bad command: {err}"))?;
210 return Ok(Incoming::Cmd(cmd));
211 }
212 let request: SendRequest =
213 serde_json::from_value(value).map_err(|_| "request is not a send JSON line".to_string())?;
214 if request.text.trim().is_empty() {
215 return Err("text is empty".to_string());
216 }
217 if request.text.len() > MAX_SEND_BYTES {
218 return Err(format!("text is longer than {MAX_SEND_BYTES} bytes"));
219 }
220 Ok(Incoming::Send(request))
221}
222
223pub fn send_cmd_via_socket(
225 sock: &Path,
226 request: &crate::CmdRequest,
227 wait: Duration,
228) -> std::io::Result<crate::CmdReply> {
229 let mut stream = SendStream::connect(sock)?;
230 stream.set_read_timeout(Some(wait))?;
231 stream.set_write_timeout(Some(Duration::from_secs(5)))?;
232 let mut line = serde_json::to_vec(request)?;
233 line.push(b'\n');
234 stream.write_all(&line)?;
235 stream.flush()?;
236 let mut reader = BufReader::new(stream);
237 let mut answer = String::new();
238 reader.read_line(&mut answer)?;
239 if answer.trim().is_empty() {
240 return Err(std::io::Error::new(std::io::ErrorKind::UnexpectedEof, "client closed the socket without an answer"));
241 }
242 serde_json::from_str(answer.trim()).map_err(|err| std::io::Error::new(std::io::ErrorKind::InvalidData, err))
243}
244
245pub(crate) fn write_cmd_reply(stream: &mut SendStream, reply: &crate::CmdReply) {
247 if let Ok(mut line) = serde_json::to_vec(reply) {
248 line.push(b'\n');
249 let _ = stream.set_write_timeout(Some(Duration::from_secs(3)));
250 let _ = stream.write_all(&line);
251 let _ = stream.flush();
252 }
253}
254
255#[cfg_attr(not(feature = "wake-grok"), allow(dead_code))]
257pub(crate) fn read_request(stream: &mut SendStream) -> Result<SendRequest, String> {
258 stream
259 .set_nonblocking(false)
260 .map_err(|err| err.to_string())?;
261 stream
262 .set_read_timeout(Some(Duration::from_secs(3)))
263 .map_err(|err| err.to_string())?;
264 let mut buf = Vec::new();
265 let mut limited = stream.take((MAX_SEND_BYTES * 2 + 1024) as u64);
266 let mut reader = BufReader::new(&mut limited);
267 reader
268 .read_until(b'\n', &mut buf)
269 .map_err(|err| err.to_string())?;
270 let request: SendRequest =
271 serde_json::from_slice(&buf).map_err(|_| "request is not a send JSON line".to_string())?;
272 if request.text.trim().is_empty() {
273 return Err("text is empty".to_string());
274 }
275 if request.text.len() > MAX_SEND_BYTES {
276 return Err(format!("text is longer than {MAX_SEND_BYTES} bytes"));
277 }
278 Ok(request)
279}
280
281pub(crate) fn write_reply(stream: &mut SendStream, reply: &SendReply) {
283 if let Ok(mut line) = serde_json::to_vec(reply) {
284 line.push(b'\n');
285 let _ = stream.set_write_timeout(Some(Duration::from_secs(3)));
286 let _ = stream.write_all(&line);
287 let _ = stream.flush();
288 }
289}
290
291#[cfg(test)]
292mod tests {
293 use std::ffi::OsStr;
294
295 use super::*;
296
297 #[test]
298 fn env_lines_skip_comments_and_strip_quotes() {
299 let parsed = parse_env_lines(
300 "# c\n\nM4A_A=1\nexport M4A_B=\"two words\"\nM4A_C='x'\nbad line\n=v\n",
301 );
302 assert_eq!(
303 parsed,
304 vec![
305 ("M4A_A".to_string(), "1".to_string()),
306 ("M4A_B".to_string(), "two words".to_string()),
307 ("M4A_C".to_string(), "x".to_string()),
308 ]
309 );
310 }
311
312 #[test]
313 fn home_dir_prefers_home_then_userprofile() {
314 assert_eq!(
315 home_dir_from(
316 Some(OsStr::new("/from-home")),
317 Some(OsStr::new("/from-profile"))
318 ),
319 Some(PathBuf::from("/from-home"))
320 );
321 assert_eq!(
322 home_dir_from(None, Some(OsStr::new("/from-profile"))),
323 Some(PathBuf::from("/from-profile"))
324 );
325 assert_eq!(
326 home_dir_from(Some(OsStr::new("")), Some(OsStr::new("/from-profile"))),
327 Some(PathBuf::from("/from-profile"))
328 );
329 assert_eq!(home_dir_from(None, None), None);
330 assert_eq!(
331 home_dir_from(Some(OsStr::new("")), Some(OsStr::new(""))),
332 None
333 );
334 }
335
336 #[test]
337 fn request_uses_as_on_the_wire() {
338 let request = SendRequest {
339 as_nick: "alice".to_string(),
340 to: "privet-mir".to_string(),
341 text: "hi".to_string(),
342 };
343 let json = serde_json::to_value(&request).expect("json");
344 assert_eq!(json["as"], "alice");
345 assert_eq!(json["to"], "privet-mir");
346 }
347}