Skip to main content

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}