1#![cfg_attr(not(ipc), allow(unused))]
2
3use std::time::Duration;
6
7use crate::{AnyResult, Event, Request, Response};
8
9use zng_task::channel::{self, ChannelError, IpcReceiver, IpcSender};
10use zng_task::parking_lot::Mutex;
11use zng_txt::Txt;
12
13type AppInitMsg = (
14 channel::IpcReceiver<Request>,
15 channel::IpcSender<Response>,
16 channel::IpcSender<Event>,
17);
18
19pub(crate) struct AppInit {
21 init_sender: channel::NamedIpcSender<AppInitMsg>,
22}
23impl AppInit {
24 pub fn new() -> Self {
25 AppInit {
26 init_sender: channel::NamedIpcSender::new().expect("failed to create init channel"),
27 }
28 }
29
30 pub fn name(&self) -> &str {
32 self.init_sender.name()
33 }
34
35 pub fn connect(self) -> AnyResult<(RequestSender, ResponseReceiver, EventReceiver)> {
37 zng_app_context::LocalContext::new().with_context(|| self.connect_blocking())
42 }
43 fn connect_blocking(self) -> AnyResult<(RequestSender, ResponseReceiver, EventReceiver)> {
44 let mut init_sender = self
45 .init_sender
46 .connect_deadline_blocking(std::time::Duration::from_secs(crate::view_timeout()))?;
47
48 let (req_sender, req_recv) = channel::ipc_unbounded()?;
49 let (rsp_sender, rsp_recv) = channel::ipc_unbounded()?;
50 let (evt_sender, evt_recv) = channel::ipc_unbounded()?;
51 init_sender.send_blocking((req_recv, rsp_sender, evt_sender))?;
52 Ok((
53 RequestSender(Mutex::new(req_sender)),
54 ResponseReceiver(Mutex::new(rsp_recv)),
55 EventReceiver(Mutex::new(evt_recv)),
56 ))
57 }
58}
59
60pub fn connect_view_process(ipc_sender_name: Txt) -> Result<ViewChannels, channel::ChannelError> {
62 let _s = tracing::trace_span!("connect_view_process").entered();
63
64 let mut init_recv = channel::IpcReceiver::<AppInitMsg>::connect(ipc_sender_name)?;
65
66 let (req_recv, rsp_sender, evt_sender) = init_recv.recv_deadline_blocking(std::time::Duration::from_secs(crate::view_timeout()))?;
67
68 Ok(ViewChannels {
69 request_receiver: RequestReceiver(Mutex::new(req_recv)),
70 response_sender: ResponseSender(Mutex::new(rsp_sender)),
71 event_sender: EventSender(Mutex::new(evt_sender)),
72 })
73}
74
75pub struct ViewChannels {
77 pub request_receiver: RequestReceiver,
82
83 pub response_sender: ResponseSender,
85
86 pub event_sender: EventSender,
88}
89
90type IpcResult<T> = Result<T, ChannelError>;
91
92pub(crate) struct RequestSender(Mutex<IpcSender<Request>>);
93impl RequestSender {
94 pub fn send(&mut self, req: Request) -> IpcResult<()> {
95 let r = self.0.get_mut().send_blocking(req);
96 if let Err(e) = &r {
97 tracing::debug!("request sender error, {e}");
98 }
99 r
100 }
101}
102impl Drop for RequestSender {
103 fn drop(&mut self) {
104 tracing::trace!("dropped RequestSender");
105 }
106}
107
108pub struct RequestReceiver(Mutex<IpcReceiver<Request>>); impl RequestReceiver {
116 pub fn recv(&mut self) -> IpcResult<Request> {
118 let r = self.0.get_mut().recv_blocking();
119 if let Err(e) = &r {
120 tracing::debug!("request receiver error, {e}");
121 }
122 r
123 }
124}
125impl Drop for RequestReceiver {
126 fn drop(&mut self) {
127 tracing::trace!("dropped RequestReceiver");
128 }
129}
130
131pub struct ResponseSender(Mutex<IpcSender<Response>>); impl ResponseSender {
140 pub fn send(&mut self, rsp: Response) -> IpcResult<()> {
148 assert!(rsp.must_be_send());
149 let r = self.0.get_mut().send_blocking(rsp);
150 if let Err(e) = &r {
151 tracing::debug!("response sender error, {e}");
152 }
153 r
154 }
155}
156impl Drop for ResponseSender {
157 fn drop(&mut self) {
158 tracing::trace!("dropped ResponseSender");
159 }
160}
161
162pub(crate) struct ResponseReceiver(Mutex<IpcReceiver<Response>>);
163impl ResponseReceiver {
164 pub fn recv(&mut self) -> IpcResult<Response> {
165 let r = self.0.get_mut().recv_blocking();
166 if let Err(e) = &r {
167 tracing::debug!("response receiver error, {e}");
168 }
169 r
170 }
171}
172impl Drop for ResponseReceiver {
173 fn drop(&mut self) {
174 tracing::trace!("dropped ResponseReceiver");
175 }
176}
177
178pub struct EventSender(Mutex<IpcSender<Event>>);
185impl EventSender {
186 pub fn send(&mut self, ev: Event) -> IpcResult<()> {
188 let r = self.0.get_mut().send_blocking(ev);
189 if let Err(e) = &r {
190 tracing::debug!("event sender error, {e}");
191 }
192 r
193 }
194}
195pub(crate) struct EventReceiver(Mutex<IpcReceiver<Event>>);
196impl EventReceiver {
197 pub fn recv(&mut self) -> IpcResult<Event> {
198 let r = self.0.get_mut().recv_blocking();
199 if let Err(e) = &r {
200 tracing::debug!("event receiver error, {e}");
201 }
202 r
203 }
204
205 pub fn recv_timeout(&mut self, duration: Duration) -> IpcResult<Event> {
206 let r = self.0.get_mut().recv_deadline_blocking(duration);
207 if let Err(e) = &r {
208 match e {
209 ChannelError::Timeout => {}
210 e => tracing::debug!("event receiver error, {e}"),
211 }
212 }
213 r
214 }
215}