ironfix_engine/application.rs
1/******************************************************************************
2 Author: Joaquín Béjar García
3 Email: jb@taunais.com
4 Date: 27/1/26
5******************************************************************************/
6
7//! Application callback interface.
8//!
9//! This module defines the callback interface for handling FIX messages,
10//! following the QuickFIX pattern with async support.
11
12use crate::outbound::OutboundMessage;
13use async_trait::async_trait;
14use ironfix_core::message::RawMessage;
15
16/// Session identifier.
17#[derive(Debug, Clone, PartialEq, Eq, Hash)]
18pub struct SessionId {
19 /// BeginString (FIX version).
20 pub begin_string: String,
21 /// Sender CompID.
22 pub sender_comp_id: String,
23 /// Target CompID.
24 pub target_comp_id: String,
25 /// Optional sender sub ID.
26 pub sender_sub_id: Option<String>,
27 /// Optional target sub ID.
28 pub target_sub_id: Option<String>,
29}
30
31impl SessionId {
32 /// Creates a new session ID.
33 #[must_use]
34 pub fn new(
35 begin_string: impl Into<String>,
36 sender_comp_id: impl Into<String>,
37 target_comp_id: impl Into<String>,
38 ) -> Self {
39 Self {
40 begin_string: begin_string.into(),
41 sender_comp_id: sender_comp_id.into(),
42 target_comp_id: target_comp_id.into(),
43 sender_sub_id: None,
44 target_sub_id: None,
45 }
46 }
47
48 /// Sets the sender sub ID.
49 #[must_use]
50 pub fn with_sender_sub_id(mut self, sub_id: impl Into<String>) -> Self {
51 self.sender_sub_id = Some(sub_id.into());
52 self
53 }
54
55 /// Sets the target sub ID.
56 #[must_use]
57 pub fn with_target_sub_id(mut self, sub_id: impl Into<String>) -> Self {
58 self.target_sub_id = Some(sub_id.into());
59 self
60 }
61}
62
63impl std::fmt::Display for SessionId {
64 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
65 write!(
66 f,
67 "{}:{}->{}",
68 self.begin_string, self.sender_comp_id, self.target_comp_id
69 )
70 }
71}
72
73/// Reason for rejecting a message.
74#[derive(Debug, Clone)]
75pub struct RejectReason {
76 /// Rejection reason code.
77 pub code: u32,
78 /// Human-readable rejection text.
79 pub text: String,
80 /// Reference tag that caused the rejection.
81 pub ref_tag: Option<u32>,
82}
83
84impl RejectReason {
85 /// Creates a new rejection reason.
86 #[must_use]
87 pub fn new(code: u32, text: impl Into<String>) -> Self {
88 Self {
89 code,
90 text: text.into(),
91 ref_tag: None,
92 }
93 }
94
95 /// Sets the reference tag.
96 #[must_use]
97 pub const fn with_ref_tag(mut self, tag: u32) -> Self {
98 self.ref_tag = Some(tag);
99 self
100 }
101}
102
103/// Application callback interface for handling FIX messages.
104///
105/// Implement this trait to receive callbacks for session events
106/// and message processing.
107///
108/// # Where each callback runs
109///
110/// [`Application::to_admin`] and [`Application::to_app`] run **inline in the
111/// session reactor**, on the message the engine is about to encode and
112/// number. That is what makes their mutations effective, and it is why they
113/// cannot be moved off the reactor: the message they hand back is the one that
114/// gets a sequence number. A slow implementation delays the session's own
115/// timers.
116///
117/// [`Application::from_admin`] also runs inline, because its verdict decides
118/// what the session does next — whether a `SequenceReset` is applied, whether
119/// the Logon is accepted — and that decision cannot be taken before the
120/// verdict exists.
121///
122/// [`Application::from_app`] runs on a **separate dispatcher task** fed by a
123/// bounded queue, so a slow handler no longer stalls socket reads, heartbeat
124/// generation or timeout detection. Messages reach it in ascending
125/// `MsgSeqNum` order. The only consequence of a rejection — a session-level
126/// Reject carrying `RefSeqNum` — is emitted by the reactor when it next polls,
127/// so it may follow messages that arrived after the rejected one. See
128/// [`Initiator::with_app_queue_capacity`](crate::Initiator::with_app_queue_capacity)
129/// for the queue depth and what happens when it fills.
130#[async_trait]
131pub trait Application: Send + Sync {
132 /// Called when a session is created.
133 ///
134 /// # Arguments
135 /// * `session_id` - The session identifier
136 async fn on_create(&self, session_id: &SessionId);
137
138 /// Called on successful logon.
139 ///
140 /// # Arguments
141 /// * `session_id` - The session identifier
142 async fn on_logon(&self, session_id: &SessionId);
143
144 /// Called on logout.
145 ///
146 /// # Arguments
147 /// * `session_id` - The session identifier
148 async fn on_logout(&self, session_id: &SessionId);
149
150 /// Called before an administrative message (Logon, Heartbeat, Logout, …)
151 /// is encoded.
152 ///
153 /// The body is mutable and the mutation is effective: the engine encodes
154 /// exactly the message this callback hands back. Fields appended here land
155 /// after the ones the session layer wrote, in insertion order. The
156 /// canonical use is stamping `Username` (553) and `Password` (554) onto the
157 /// outbound Logon.
158 ///
159 /// The standard header is not visible: it is stamped after this returns,
160 /// and `MsgSeqNum` (34) is not decided until the frame exists. A body that
161 /// carries a tag the engine stamps itself, or a value with no legal wire
162 /// form, is refused — the message is dropped with a warning and the session
163 /// continues, having spent no sequence number. See
164 /// [`crate::outbound::RESERVED_TAGS`].
165 ///
166 /// # Arguments
167 /// * `message` - The message about to be encoded (mutable)
168 /// * `session_id` - The session identifier
169 async fn to_admin(&self, message: &mut OutboundMessage, session_id: &SessionId);
170
171 /// Called when an admin message is received.
172 ///
173 /// # Arguments
174 /// * `message` - The received message
175 /// * `session_id` - The session identifier
176 ///
177 /// # Returns
178 /// `Ok(())` to accept, `Err(RejectReason)` to reject.
179 #[allow(clippy::wrong_self_convention)]
180 async fn from_admin(
181 &self,
182 message: &RawMessage<'_>,
183 session_id: &SessionId,
184 ) -> Result<(), RejectReason>;
185
186 /// Called before an application message is encoded.
187 ///
188 /// Same contract as [`Application::to_admin`]: the body is mutable, the
189 /// mutation is effective, and the header — including `MsgSeqNum` (34) — is
190 /// stamped afterwards.
191 ///
192 /// # Arguments
193 /// * `message` - The message about to be encoded (mutable)
194 /// * `session_id` - The session identifier
195 async fn to_app(&self, message: &mut OutboundMessage, session_id: &SessionId);
196
197 /// Called when an application message is received.
198 ///
199 /// # Arguments
200 /// * `message` - The received message
201 /// * `session_id` - The session identifier
202 ///
203 /// # Returns
204 /// `Ok(())` to accept, `Err(RejectReason)` to reject.
205 #[allow(clippy::wrong_self_convention)]
206 async fn from_app(
207 &self,
208 message: &RawMessage<'_>,
209 session_id: &SessionId,
210 ) -> Result<(), RejectReason>;
211}
212
213/// Default no-op application implementation.
214#[derive(Debug, Default)]
215pub struct NoOpApplication;
216
217#[async_trait]
218impl Application for NoOpApplication {
219 async fn on_create(&self, _session_id: &SessionId) {}
220
221 async fn on_logon(&self, _session_id: &SessionId) {}
222
223 async fn on_logout(&self, _session_id: &SessionId) {}
224
225 async fn to_admin(&self, _message: &mut OutboundMessage, _session_id: &SessionId) {}
226
227 async fn from_admin(
228 &self,
229 _message: &RawMessage<'_>,
230 _session_id: &SessionId,
231 ) -> Result<(), RejectReason> {
232 Ok(())
233 }
234
235 async fn to_app(&self, _message: &mut OutboundMessage, _session_id: &SessionId) {}
236
237 async fn from_app(
238 &self,
239 _message: &RawMessage<'_>,
240 _session_id: &SessionId,
241 ) -> Result<(), RejectReason> {
242 Ok(())
243 }
244}
245
246#[cfg(test)]
247mod tests {
248 use super::*;
249
250 #[test]
251 fn test_session_id() {
252 let id = SessionId::new("FIX.4.4", "SENDER", "TARGET");
253 assert_eq!(id.begin_string, "FIX.4.4");
254 assert_eq!(id.sender_comp_id, "SENDER");
255 assert_eq!(id.target_comp_id, "TARGET");
256 assert_eq!(id.to_string(), "FIX.4.4:SENDER->TARGET");
257 }
258
259 #[test]
260 fn test_reject_reason() {
261 let reason = RejectReason::new(1, "Invalid tag").with_ref_tag(35);
262 assert_eq!(reason.code, 1);
263 assert_eq!(reason.text, "Invalid tag");
264 assert_eq!(reason.ref_tag, Some(35));
265 }
266
267 #[tokio::test]
268 async fn test_noop_application() {
269 let app = NoOpApplication;
270 let session_id = SessionId::new("FIX.4.4", "SENDER", "TARGET");
271
272 app.on_create(&session_id).await;
273 app.on_logon(&session_id).await;
274 app.on_logout(&session_id).await;
275 }
276}