boxmux_lib/socket_loop.rs
1use crate::model::common::{run_socket_function, SocketFunction};
2use crate::thread_manager::Runnable;
3use crate::{AppContext, FieldUpdate};
4use std::fs;
5use std::io::{Read, Write};
6use std::os::unix::net::UnixListener;
7use std::sync::mpsc;
8
9use crate::thread_manager::*;
10
11use uuid::Uuid;
12
13create_runnable!(
14 SocketLoop,
15 |_inner: &mut RunnableImpl, _app_context: AppContext, _messages: Vec<Message>| -> bool { true },
16 |inner: &mut RunnableImpl,
17 app_context: AppContext,
18 _messages: Vec<Message>|
19 -> (bool, AppContext) {
20 let socket_path = "/tmp/boxmux.sock";
21 // Remove the stale socket file if it exists
22 if std::path::Path::new(socket_path).exists() {
23 let _ = fs::remove_file(socket_path);
24 }
25
26 let listener = match UnixListener::bind(socket_path) {
27 Ok(listener) => {
28 log::info!("Listening on socket: {}", socket_path);
29 listener
30 }
31 Err(err) => {
32 log::error!("Failed to bind to socket {}: {}", socket_path, err);
33 return (false, app_context);
34 }
35 };
36
37 for stream in listener.incoming() {
38 match stream {
39 Ok(mut stream) => {
40 let mut buffer = String::new();
41 match stream.read_to_string(&mut buffer) {
42 Ok(_size) => {
43 let trimmed_message = buffer.trim();
44 log::debug!("Received socket message: {}", trimmed_message);
45
46 // Parse JSON message as SocketFunction and execute directly
47 if !trimmed_message.is_empty() {
48 match serde_json::from_str::<SocketFunction>(trimmed_message) {
49 Ok(socket_function) => {
50 log::debug!(
51 "Parsed socket function: {:?}",
52 socket_function
53 );
54
55 // Execute socket function and send resulting messages
56 match run_socket_function(socket_function, &app_context) {
57 Ok((_updated_context, messages)) => {
58 // Update app_context if it was modified
59 // Note: app_context is typically not modified by socket functions
60 // but we maintain the pattern for consistency
61
62 // Send all resulting messages to the thread manager
63 for message in messages {
64 inner.send_message(message);
65 }
66
67 // Send success acknowledgment
68 if let Err(err) = stream.write_all(
69 b"Socket function executed successfully.",
70 ) {
71 log::error!(
72 "Error sending success response: {}",
73 err
74 );
75 }
76 }
77 Err(err) => {
78 let error_msg = format!(
79 "Socket function execution failed: {}",
80 err
81 );
82 log::error!("{}", error_msg);
83
84 // Send error response to client
85 if let Err(write_err) =
86 stream.write_all(error_msg.as_bytes())
87 {
88 log::error!(
89 "Error sending error response: {}",
90 write_err
91 );
92 }
93 }
94 }
95 }
96 Err(parse_err) => {
97 let error_msg =
98 format!("Invalid socket function JSON: {}", parse_err);
99 log::error!("{}", error_msg);
100
101 // Send parse error response to client
102 if let Err(write_err) =
103 stream.write_all(error_msg.as_bytes())
104 {
105 log::error!(
106 "Error sending parse error response: {}",
107 write_err
108 );
109 }
110 }
111 }
112 }
113 }
114 Err(err) => {
115 log::error!("Error receiving message: {}", err);
116 }
117 }
118 }
119 Err(err) => {
120 log::error!("Error accepting connection: {}", err);
121 }
122 }
123 }
124
125 (true, app_context)
126 }
127);