Skip to main content

zng_view_api/
ipc.rs

1#![cfg_attr(not(ipc), allow(unused))]
2
3//! IPC types.
4
5use 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
19/// Call `new`, then spawn the view-process using the `name` then call `connect`.
20pub(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    /// Unique name for the view-process to find this channel.
31    pub fn name(&self) -> &str {
32        self.init_sender.name()
33    }
34
35    /// Tries to connect to the view-process and receive the actual channels.
36    pub fn connect(self) -> AnyResult<(RequestSender, ResponseReceiver, EventReceiver)> {
37        // avoid app context because `connect_deadline_blocking` uses `block_on`
38        // and that logs a warning. This is not async because it is tricky to
39        // await when the app loop is just starting, and it will barely block
40        // when the view-process spawns correctly
41        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
60/// Start the view-process server and waits for `(request, response, event)`.
61pub 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
75/// Channels that must be used for implementing a view-process.
76pub struct ViewChannels {
77    /// View implementers must receive requests from this channel, call [`Api::respond`] and then
78    /// return the response using the `response_sender`.
79    ///
80    /// [`Api::respond`]: crate::Api::respond
81    pub request_receiver: RequestReceiver,
82
83    /// View implementers must synchronously send one response per request received in `request_receiver`.
84    pub response_sender: ResponseSender,
85
86    /// View implements must send events using this channel. Events can be asynchronous.
87    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
108/// Requests channel end-point.
109///
110/// View-process implementers must receive [`Request`], call [`Api::respond`] and then use a [`ResponseSender`]
111/// to send back the response.
112///
113/// [`Api::respond`]: crate::Api::respond
114pub struct RequestReceiver(Mutex<IpcReceiver<Request>>); // Mutex for Sync
115impl RequestReceiver {
116    /// Receive one [`Request`].
117    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
131/// Responses channel entry-point.
132///
133/// View-process implementers must send [`Response`] returned by [`Api::respond`] using this sender.
134///
135/// Requests are received using [`RequestReceiver`] a response must be send for each request, synchronously.
136///
137/// [`Api::respond`]: crate::Api::respond
138pub struct ResponseSender(Mutex<IpcSender<Response>>); // Mutex for Sync
139impl ResponseSender {
140    /// Send a response.
141    ///
142    /// # Panics
143    ///
144    /// If the `rsp` is not [`must_be_send`].
145    ///
146    /// [`must_be_send`]: Response::must_be_send
147    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
178/// Event channel entry-point.
179///
180/// View-process implementers must send [`Event`] messages using this sender. The events
181/// can be asynchronous, not related to the [`Api::respond`] calls.
182///
183/// [`Api::respond`]: crate::Api::respond
184pub struct EventSender(Mutex<IpcSender<Event>>);
185impl EventSender {
186    /// Send an event notification.
187    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}