Skip to main content

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}