nerve_ipc_core/dispatch.rs
1//! Transport-agnostic frame dispatch.
2//!
3//! [`dispatch_frame`] is the single function both transport layers call for
4//! every received NERVE frame. It reads the decoded frame, updates the
5//! [`crate::RequestTable`] as needed, and returns a [`DispatchAction`] telling
6//! the transport what to do next.
7//!
8//! `dispatch_frame` performs no I/O. All socket reads and writes are the
9//! responsibility of the calling transport layer ([`crate::server`] or
10//! [`crate::ws_server`]).
11//!
12//! # AI daemon boundary
13//!
14//! When a `SearchQuery` frame is received, `dispatch_frame` registers the
15//! request in the [`crate::RequestTable`] and returns
16//! [`DispatchAction::ForwardToAiDaemon`]. The AI daemon (a future milestone)
17//! is responsible for consuming the query, checking cancellation, and streaming
18//! results back. The transport loop must not send any reply itself for this
19//! variant.
20
21use nerve_protocol::codec::encode;
22use nerve_protocol::types::{FrameFlags, MessageType, RequestId};
23use nerve_protocol::{Frame, ProtocolError};
24
25use crate::request_table::RequestTable;
26
27/// Outcome returned to the connection loop after dispatching a frame.
28///
29/// Marked `#[non_exhaustive]` because new variants will be added as the AI
30/// daemon layer is built out. Downstream code must include a wildcard arm
31/// when matching on this type.
32///
33/// # Examples
34///
35/// ```
36/// use nerve_ipc_core::DispatchAction;
37///
38/// fn handle(action: DispatchAction) -> Option<Vec<u8>> {
39/// match action {
40/// DispatchAction::Handled => None,
41/// DispatchAction::Reply(bytes) => Some(bytes),
42/// DispatchAction::ForwardToAiDaemon(_) => None,
43/// // Required: new variants will be added as the AI daemon is built out.
44/// _ => None,
45/// }
46/// }
47/// ```
48#[non_exhaustive]
49pub enum DispatchAction {
50 /// Frame was handled; no response needs to be sent to the client.
51 Handled,
52 /// Frame generated a response that the transport must write back to the
53 /// client. The bytes are a complete, encoded NERVE frame.
54 Reply(Vec<u8>),
55 /// A [`MessageType::SearchQuery`] was received. The request has been
56 /// registered in the [`RequestTable`] so that a subsequent
57 /// [`MessageType::Cancel`] can mark it cancelled. The AI daemon (a future
58 /// milestone) is responsible for consuming and responding to this request.
59 /// The enclosed [`RequestId`] is the ID already stored in the request table.
60 ForwardToAiDaemon(RequestId),
61}
62
63/// Dispatch a single decoded frame and return the action to take.
64///
65/// This function is transport-agnostic: it does not perform any I/O.
66/// The caller is responsible for sending any [`DispatchAction::Reply`] bytes
67/// back to the client, whether over a Unix socket or a WebSocket.
68///
69/// Must be fast, non-blocking, and deterministic.
70///
71/// # Errors
72///
73/// Returns [`ProtocolError`] only when encoding the `Ping` reply fails. All
74/// other message types either produce no reply or register state in the
75/// [`RequestTable`] without encoding a response.
76///
77/// # Examples
78///
79/// ```
80/// use nerve_ipc_core::{dispatch_frame, DispatchAction, RequestTable};
81/// use nerve_protocol::codec::{decode, encode};
82/// use nerve_protocol::types::{FrameFlags, MessageType, RequestId};
83///
84/// let bytes = encode(MessageType::Ping, FrameFlags::empty(), RequestId(1), b"").unwrap();
85/// let frame = decode(&bytes).unwrap();
86/// let mut table = RequestTable::new();
87///
88/// match dispatch_frame(frame, &mut table).unwrap() {
89/// DispatchAction::Reply(_reply) => { /* send reply to the client */ }
90/// DispatchAction::Handled => {}
91/// DispatchAction::ForwardToAiDaemon(_) => {}
92/// _ => {}
93/// }
94/// ```
95pub fn dispatch_frame(
96 frame: Frame<'_>,
97 requests: &mut RequestTable,
98) -> Result<DispatchAction, ProtocolError> {
99 let msg_type = match MessageType::try_from(frame.header.msg_type) {
100 Ok(t) => t,
101 Err(_) => return Ok(DispatchAction::Handled),
102 };
103
104 match msg_type {
105 MessageType::Ping => {
106 let bytes = encode_ping_reply(frame)?;
107 Ok(DispatchAction::Reply(bytes))
108 }
109
110 MessageType::Cancel => {
111 handle_cancel(frame, requests);
112 Ok(DispatchAction::Handled)
113 }
114
115 // SearchQuery belongs to the AI daemon layer. Register the request
116 // here so cancellation works, then let the caller forward it.
117 MessageType::SearchQuery => {
118 let req_id = RequestId(frame.header.request_id);
119 requests.insert(req_id);
120 Ok(DispatchAction::ForwardToAiDaemon(req_id))
121 }
122
123 MessageType::AgentTaskStart => {
124 let req_id = RequestId(frame.header.request_id);
125 requests.insert(req_id);
126 Ok(DispatchAction::Handled)
127 }
128
129 MessageType::AgentTaskEvent => Ok(DispatchAction::Handled),
130
131 MessageType::AgentTaskDone => {
132 let req_id = RequestId(frame.header.request_id);
133 requests.remove(req_id);
134 Ok(DispatchAction::Handled)
135 }
136
137 // AiToken / SearchResult flow from the AI daemon to the client.
138 // Ignore if received in this direction.
139 _ => Ok(DispatchAction::Handled),
140 }
141}
142
143fn encode_ping_reply(frame: Frame<'_>) -> Result<Vec<u8>, ProtocolError> {
144 encode(
145 MessageType::Ping,
146 FrameFlags::FINAL,
147 RequestId(frame.header.request_id),
148 frame.payload,
149 )
150}
151
152fn handle_cancel(frame: Frame<'_>, requests: &mut RequestTable) {
153 let req_id = RequestId(frame.header.request_id);
154 requests.cancel(req_id);
155}